use std::{
collections::BTreeMap,
error::Error as StdError,
future::Future,
panic::{AssertUnwindSafe, catch_unwind},
path::{Path, PathBuf},
pin::Pin,
sync::Arc,
task::{Context, Poll},
};
use clap::error::ErrorKind;
use markdown_compiler::{
ContentCandidateStore, ContentCandidateStoreError, ContentTreeDigest, ContentTreeLimits,
ContentValidationErrors, DiscoveredContentTree, PostId, PostRevisionDigest,
PrepareContentError, PreparedContent, SiteSnapshotDigest, discover_content_tree,
prepare_content,
};
use time::OffsetDateTime;
use tokio::task::{JoinError, JoinHandle, JoinSet};
use tokio_util::sync::CancellationToken;
use tracing::Instrument as _;
use crate::{
admin::{
AdminSecurityState, AdminServer, AdminSessionPolicy, origin::AdminBind,
runtime_admin_router,
},
backup_health::BackupHealth,
cli::{ServerInvocation, parse_process_invocation},
config::{HostConfiguration, HostConfigurationLoader, SourceConfigurationView},
content_sync::ContentSync,
database::{self, DatabaseStore},
domain::{
auth::{Argon2idPolicy, store::ConfiguredLoginProviders},
mail::{
config::MailConfiguration,
dispatch::MailDispatcher,
feedback::FeedbackWorker,
retention::MailRetention,
runtime::{PreparedMail, prepare_mail},
ui::MailUiState,
},
profile::TipRecipientProjection,
publication::{
PublicLedgerProjection, SourceCommit,
activation::{
PreparedPublicationRecovery, PublicationActivationError, PublicationCoordinator,
PublicationCoordinatorActor, PublicationCoordinatorHandle, observed_post_revisions,
},
scheduler::PublicationScheduler,
store::{
InstallStartupSnapshot, ObservedPostRevision, RecoverablePublicationActivation,
RetainedReleaseInput, StartupSnapshotState,
},
},
},
error::{
ApplicationError, CriticalTaskName, ProcessError, ProcessExit, ShutdownSignal, StartupStage,
},
frontend_assets::{FrontendAssetManifest, embedded_manifest},
git_sync::GitSync,
identity_bootstrap,
metrics::{Metrics, MetricsCollector, MetricsServer},
observability::{initialize_logging, task_span},
process_lock::{ProcessLock, ProcessLockError},
render::{
CatalogBuildError, CatalogRetentionError, ContentCatalog, ContentCompiler, SiteSnapshot,
render_site_shell, snapshot_store,
},
restore,
source_bootstrap::{configure_source, generate_source_key},
source_provenance::{SourceCommitDiscovery, discover_source_commit},
source_sync::{ManagedSourceEngine, SourceSyncHandle},
web::{PublicServer, PublicState, Readiness, public_router_with_routes},
};
#[cfg(test)]
use crate::{domain::publication::PublishedPostRevision, render::compile_content_catalog};
type ShutdownFuture = Pin<Box<dyn Future<Output = Result<(), ApplicationError>> + Send>>;
type CriticalTaskFuture = Pin<Box<dyn Future<Output = CriticalTaskResult> + Send>>;
type CriticalTaskResult = Result<(), CriticalTaskFailure>;
type CriticalTaskFailure = Box<dyn std::error::Error + Send + Sync>;
const PUBLICATION_COORDINATOR_QUEUE_CAPACITY: usize = 32;
pub(crate) struct Application {
_process_lock: ProcessLock,
_database: DatabaseStore,
publication_coordinator: PublicationCoordinatorHandle,
runtime: ApplicationRuntime,
#[cfg(test)]
public_addr: std::net::SocketAddr,
#[cfg(test)]
admin_addr: std::net::SocketAddr,
}
struct ApplicationRuntime {
readiness: Readiness,
cancellation: CancellationToken,
database_shutdown: CancellationToken,
shutdown: ShutdownFuture,
critical_tasks: JoinSet<(CriticalTaskName, CriticalTaskCompletion)>,
database_writer: Option<JoinHandle<(CriticalTaskName, CriticalTaskCompletion)>>,
}
struct StartupConfiguration {
_process_lock: ProcessLock,
_host: HostConfiguration,
_content_tree: DiscoveredContentTree,
content: PreparedContent,
}
struct StartupHostConfiguration {
process_lock: ProcessLock,
host: HostConfiguration,
}
struct CompiledStartupContent {
catalog: Arc<ContentCatalog>,
observed_posts: Vec<ObservedPostRevision>,
source_commit: Option<SourceCommit>,
content_digest: ContentTreeDigest,
}
struct ServingState {
readiness: Readiness,
publication_coordinator: PublicationCoordinatorHandle,
publication_actor: PublicationCoordinatorActor,
public_server: PublicServer,
admin_server: AdminServer,
metrics_server: MetricsServer,
metrics_collector: MetricsCollector,
mail_dispatcher: Option<MailDispatcher>,
mail_feedback: Option<FeedbackWorker>,
}
struct StartedDatabase {
store: DatabaseStore,
shutdown: CancellationToken,
task: JoinHandle<(CriticalTaskName, CriticalTaskCompletion)>,
security: AdminSecurityState,
}
impl StartupConfiguration {
#[cfg(test)]
fn load_with_discovery<Discover>(
config_path: PathBuf,
discover: Discover,
) -> Result<Self, ProcessError>
where
Discover: FnOnce(
&Path,
ContentTreeLimits,
) -> Result<DiscoveredContentTree, ContentValidationErrors>,
{
StartupHostConfiguration::load(config_path)?.discover_with(discover)
}
}
impl StartupHostConfiguration {
fn load(config_path: PathBuf) -> Result<Self, ProcessError> {
let host = HostConfigurationLoader::from_process_working_directory()?.load(&config_path)?;
let host_view = host.view();
let process_lock = match ProcessLock::acquire(host_view.runtime_root) {
Ok(process_lock) => process_lock,
Err(ProcessLockError::AlreadyRunning) => return Err(ProcessError::AlreadyRunning),
Err(error) => {
return Err(startup_failure(
StartupStage::ProcessLock,
"acquire the process lock",
error,
));
}
};
Ok(Self { process_lock, host })
}
fn discover_with<Discover>(
self,
discover: Discover,
) -> Result<StartupConfiguration, ProcessError>
where
Discover: FnOnce(
&Path,
ContentTreeLimits,
) -> Result<DiscoveredContentTree, ContentValidationErrors>,
{
let host_view = self.host.view();
let content_tree = discover(host_view.content_root, host_view.content_limits)?;
let content = prepare_content(&content_tree).map_err(|error| match error {
PrepareContentError::InvalidContent(source) => ProcessError::Validation(source),
PrepareContentError::AssetResolution(source) => {
startup_failure(StartupStage::Content, "resolve content assets", source)
}
})?;
Ok(StartupConfiguration {
_process_lock: self.process_lock,
_host: self.host,
_content_tree: content_tree,
content,
})
}
}
pub async fn run_until_stop() -> ProcessExit {
initialize_logging();
let invocation = match parse_process_invocation() {
Ok(invocation) => invocation,
Err(error) => return report_command_error(error),
};
let result: Result<(), ProcessError> = async {
match invocation {
ServerInvocation::CheckpointManifest {
config_path,
database_file,
plan_file,
ltx_root,
artifact_root,
output,
} => {
restore::checkpoint::write_manifest(
config_path,
database_file,
plan_file,
ltx_root,
artifact_root,
output,
)
.await
}
ServerInvocation::VerifyCheckpoint {
manifest_file,
ltx_root,
artifact_root,
} => {
restore::checkpoint::verify_checkpoint(manifest_file, ltx_root, artifact_root).await
}
ServerInvocation::RestoreReplica {
config_path,
database_file,
artifact_root,
manifest_file,
ltx_root,
} => {
restore::restore(
config_path,
database_file,
artifact_root,
restore::RestoreManifest::Replica {
path: manifest_file,
ltx_root,
},
)
.await
}
ServerInvocation::ExportBackup {
config_path,
database_file,
} => restore::export_backup(config_path, database_file).await,
ServerInvocation::Restore {
config_path,
database_file,
artifact_root,
manifest_file,
} => {
restore::restore(
config_path,
database_file,
artifact_root,
restore::RestoreManifest::Bundle(manifest_file),
)
.await
}
ServerInvocation::Serve { config_path } => {
let startup = StartupHostConfiguration::load(config_path)?;
let application = match startup.host.view().source {
SourceConfigurationView::ExternalCheckout => {
Application::build(startup.discover_with(discover_content_tree)?).await?
}
SourceConfigurationView::ManagedGit { .. } => {
Application::build_managed(startup).await?
}
};
application
.run_until_stop()
.await
.map_err(ProcessError::from)
}
ServerInvocation::BootstrapIdentity {
config_path,
credential,
} => identity_bootstrap::bootstrap_owner(config_path, credential).await,
ServerInvocation::ConfigureSource {
config_path,
request,
idempotency_key,
} => configure_source(config_path, request, idempotency_key).await,
ServerInvocation::GenerateSourceKey {
config_path,
private_key_file,
} => generate_source_key(config_path, private_key_file).await,
}
}
.await;
match result {
Ok(()) => ProcessExit::Success,
Err(error) => {
let exit = error.exit();
tracing::error!(
error = %error,
category = error.category(),
exit_code = exit.code(),
"server process failed"
);
exit
}
}
}
fn report_command_error(error: clap::Error) -> ProcessExit {
let exit = match error.kind() {
ErrorKind::DisplayHelp | ErrorKind::DisplayVersion => ProcessExit::Success,
_ => ProcessExit::Usage,
};
if let Err(print_error) = error.print() {
tracing::error!(
error = %print_error,
"failed to print command output"
);
return ProcessExit::Internal;
}
exit
}
impl Application {
async fn build(startup: StartupConfiguration) -> Result<Self, ProcessError> {
let cancellation = CancellationToken::new();
let host = startup._host.view();
let content_root = host.content_root.to_path_buf();
let state_root = host.state_root.to_path_buf();
let content_limits = host.content_limits;
let frontend = embedded_manifest();
frontend.validate().map_err(|error| {
startup_failure(
StartupStage::FrontendAssets,
"validate embedded frontend assets",
error,
)
})?;
let shutdown = install_termination_signal()?;
let candidate_store =
ContentCandidateStore::open(&state_root, content_limits).map_err(|error| {
startup_failure(
StartupStage::Content,
"open the retained content candidate store",
error,
)
})?;
candidate_store
.retain(&startup._content_tree)
.map_err(|error| {
startup_failure(
StartupStage::Content,
"retain the startup content candidate",
error,
)
})?;
let content_compiler = ContentCompiler::discover().map_err(|error| {
startup_failure(
StartupStage::Content,
"initialize the content rendering pipeline",
error,
)
})?;
let compiled = compile_startup_content(&startup, host.content_root, &content_compiler)?;
let active_content_digest = compiled.content_digest.clone();
let database = start_database(&startup._host).await?;
let source = SourceSyncHandle::external_checkout(database.store.source.clone());
let serving_state = match prepare_serving_state(ServingStateInput {
database: &database.store,
compiled,
frontend,
public_bind: host.public_bind,
admin_bind: host.admin_bind,
metrics_bind: host.metrics_bind,
backup: BackupHealth::new(host.backup),
security: database.security.clone(),
mail: host.mail,
cancellation: cancellation.clone(),
candidate_store: &candidate_store,
content_compiler: &content_compiler,
source,
})
.await
{
Ok(setup) => setup,
Err(error) => {
return Err(close_started_database(database, error).await);
}
};
let content_sync = ContentSync::new(
content_root,
content_limits,
candidate_store,
active_content_digest,
serving_state.publication_coordinator.clone(),
cancellation.clone(),
content_compiler,
);
let content_task = CriticalTask::new(CriticalTaskName::ContentSync, content_sync.run());
Ok(Self::assemble(
startup._process_lock,
database,
serving_state,
cancellation,
shutdown,
content_task,
))
}
async fn build_managed(startup: StartupHostConfiguration) -> Result<Self, ProcessError> {
let cancellation = CancellationToken::new();
let host = startup.host.view();
let state_root = host.state_root.to_path_buf();
let content_limits = host.content_limits;
let (mirror_root, credentials, process_limits) = match host.source {
SourceConfigurationView::ManagedGit {
mirror_root,
credentials,
limits,
} => (mirror_root, credentials, limits),
SourceConfigurationView::ExternalCheckout => {
return Err(ProcessError::ManagedSourceDisabled);
}
};
let frontend = embedded_manifest();
frontend.validate().map_err(|error| {
startup_failure(
StartupStage::FrontendAssets,
"validate embedded frontend assets",
error,
)
})?;
let shutdown = install_termination_signal()?;
let database = start_database(&startup.host).await?;
match database.store.source.configuration().await {
Ok(Some(_)) => {}
Ok(None) => {
return Err(close_started_database(
database,
ProcessError::SourceConfigurationRequired,
)
.await);
}
Err(error) => {
let error = startup_failure(
StartupStage::Source,
"load durable managed-source settings",
error,
);
return Err(close_started_database(database, error).await);
}
};
let git = match GitSync::discover(mirror_root, credentials, process_limits, content_limits)
{
Ok(git) => git,
Err(error) => {
let error = startup_failure(
StartupStage::Source,
"initialize the managed Git transport",
error,
);
return Err(close_started_database(database, error).await);
}
};
let candidate_store = match ContentCandidateStore::open(&state_root, content_limits) {
Ok(store) => store,
Err(error) => {
let error = startup_failure(
StartupStage::Content,
"open the retained content candidate store",
error,
);
return Err(close_started_database(database, error).await);
}
};
let content_compiler = match ContentCompiler::discover() {
Ok(compiler) => compiler,
Err(error) => {
let error = startup_failure(
StartupStage::Content,
"initialize the content rendering pipeline",
error,
);
return Err(close_started_database(database, error).await);
}
};
let (source_engine, source) = ManagedSourceEngine::new(
database.store.source.clone(),
git,
candidate_store.clone(),
content_compiler.clone(),
cancellation.clone(),
);
let candidate = match source_engine.prepare_startup().await {
Ok(candidate) => candidate,
Err(error) => {
let error = startup_failure(
StartupStage::Source,
"prepare the managed source head",
error,
);
return Err(close_started_database(database, error).await);
}
};
let compiled = CompiledStartupContent {
observed_posts: observed_post_revisions(&candidate.catalog),
catalog: candidate.catalog,
source_commit: Some(candidate.source_commit),
content_digest: candidate.content_digest,
};
let serving_state = match prepare_serving_state(ServingStateInput {
database: &database.store,
compiled,
frontend,
public_bind: host.public_bind,
admin_bind: host.admin_bind,
metrics_bind: host.metrics_bind,
backup: BackupHealth::new(host.backup),
security: database.security.clone(),
mail: host.mail,
cancellation: cancellation.clone(),
candidate_store: &candidate_store,
content_compiler: &content_compiler,
source,
})
.await
{
Ok(setup) => setup,
Err(error) => {
return Err(close_started_database(database, error).await);
}
};
let source_sync = source_engine.into_live(serving_state.publication_coordinator.clone());
let source_task = CriticalTask::new(CriticalTaskName::SourceSync, source_sync.run());
Ok(Self::assemble(
startup.process_lock,
database,
serving_state,
cancellation,
shutdown,
source_task,
))
}
fn assemble(
process_lock: ProcessLock,
database: StartedDatabase,
serving_state: ServingState,
cancellation: CancellationToken,
shutdown: ShutdownFuture,
source_task: CriticalTask,
) -> Self {
let ServingState {
readiness,
publication_coordinator,
publication_actor,
public_server,
admin_server,
metrics_server,
metrics_collector,
mail_dispatcher,
mail_feedback,
} = serving_state;
#[cfg(test)]
let public_addr = public_server.local_addr;
#[cfg(test)]
let admin_addr = admin_server.local_addr;
let public_task = CriticalTask::new(
CriticalTaskName::PublicServer,
public_server.serve(cancellation.clone()),
);
let admin_task = CriticalTask::new(
CriticalTaskName::AdminServer,
admin_server.serve(cancellation.clone()),
);
let metrics_task = CriticalTask::new(
CriticalTaskName::MetricsServer,
metrics_server.serve(cancellation.clone()),
);
let collector_task = CriticalTask::new(
CriticalTaskName::MetricsCollector,
metrics_collector.run(cancellation.clone()),
);
let publication_actor_task = CriticalTask::new(
CriticalTaskName::PublicationCoordinator,
publication_actor.run(cancellation.clone()),
);
let scheduler = PublicationScheduler::new(
database.store.publications.clone(),
publication_coordinator.clone(),
publication_coordinator.scheduler_wakeup(),
cancellation.clone(),
);
let scheduler_task = CriticalTask::new(CriticalTaskName::Scheduler, scheduler.run());
let retention = MailRetention::new(database.store.subscribers.clone());
let retention_task = CriticalTask::new(
CriticalTaskName::MailRetention,
retention.run(cancellation.clone()),
);
let mut tasks = vec![
retention_task,
publication_actor_task,
public_task,
admin_task,
metrics_task,
collector_task,
source_task,
scheduler_task,
];
if let Some(dispatcher) = mail_dispatcher {
tasks.push(CriticalTask::new(
CriticalTaskName::MailDispatch,
dispatcher.run(cancellation.clone()),
));
}
if let Some(feedback) = mail_feedback {
tasks.push(CriticalTask::new(
CriticalTaskName::MailFeedback,
feedback.run(cancellation.clone()),
));
}
Self {
_process_lock: process_lock,
_database: database.store,
publication_coordinator,
runtime: ApplicationRuntime::with_database_writer(
readiness,
cancellation,
database.shutdown,
shutdown,
tasks,
database.task,
),
#[cfg(test)]
public_addr,
#[cfg(test)]
admin_addr,
}
}
async fn run_until_stop(self) -> Result<(), ApplicationError> {
let Self {
_process_lock: process_lock,
_database: database,
publication_coordinator,
runtime,
#[cfg(test)]
public_addr: _,
#[cfg(test)]
admin_addr: _,
} = self;
let runtime_result = runtime.run_until_stop().await;
drop(publication_coordinator);
drop(database);
drop(process_lock);
runtime_result
}
}
fn compile_startup_content(
startup: &StartupConfiguration,
content_root: &Path,
compiler: &ContentCompiler,
) -> Result<CompiledStartupContent, ProcessError> {
let content_digest = startup._content_tree.digest();
let catalog = Arc::new(compiler.compile(&startup.content).map_err(|error| {
startup_failure(StartupStage::Content, "compile the content catalog", error)
})?);
let observed_posts = observed_post_revisions(&catalog);
let source_commit = match discover_source_commit(content_root) {
SourceCommitDiscovery::Discovered(commit) => Some(commit),
SourceCommitDiscovery::Unavailable(reason) => {
tracing::warn!(?reason, "content source commit is unavailable");
None
}
};
Ok(CompiledStartupContent {
catalog,
observed_posts,
source_commit,
content_digest,
})
}
pub(crate) async fn verify_restore_content(
database: &DatabaseStore,
state_root: &Path,
limits: ContentTreeLimits,
) -> Result<(), ProcessError> {
let inputs = RestoreContentInputs::load(database).await?;
let state_root = state_root.to_path_buf();
tokio::task::spawn_blocking(move || inputs.verify(&state_root, limits))
.await
.map_err(|error| {
startup_failure(
StartupStage::Content,
"await offline restore verification",
error,
)
})?
}
struct RestoreContentInputs {
tip_recipient: Option<TipRecipientProjection>,
startup: StartupSnapshotState,
installed_content: Option<ContentTreeDigest>,
retained: Vec<RetainedReleaseInput>,
}
impl RestoreContentInputs {
async fn load(database: &DatabaseStore) -> Result<Self, ProcessError> {
let tip_recipient = database
.profiles
.effective_tip_recipient()
.await
.map_err(|error| {
startup_failure(
StartupStage::Database,
"verify restored tip projection",
error,
)
})?;
let startup = database
.publications
.startup_snapshot_state()
.await
.map_err(|error| {
startup_failure(
StartupStage::Database,
"verify restored publication ledger",
error,
)
})?;
let source = database.source.status().await.map_err(|error| {
startup_failure(
StartupStage::Database,
"verify restored source ledger",
error,
)
})?;
let retained = database
.publications
.retained_release_inputs()
.await
.map_err(|error| {
startup_failure(
StartupStage::Database,
"verify retained restore approvals",
error,
)
})?;
Ok(Self {
tip_recipient,
startup,
installed_content: source.installation.map(|source| source.content_digest),
retained,
})
}
fn verify(&self, state_root: &Path, limits: ContentTreeLimits) -> Result<(), ProcessError> {
let candidates = ContentCandidateStore::open(state_root, limits).map_err(|error| {
startup_failure(
StartupStage::Content,
"inspect restored candidate archives",
error,
)
})?;
let compiler = ContentCompiler::discover().map_err(|error| {
startup_failure(StartupStage::Content, "initialize restore compiler", error)
})?;
let catalogs = compile_retained_catalogs(&candidates, &compiler).map_err(|error| {
startup_failure(
StartupStage::Content,
"compile restored candidate archives",
error,
)
})?;
let base = catalogs.values().next().ok_or_else(|| {
startup_failure(
StartupStage::Content,
"verify retained restore inputs",
RetainedCatalogError::CandidateUnavailable,
)
})?;
let pins = self.startup.ledger.revision_keys().chain(
self.retained
.iter()
.map(|input| (input.post_id.clone(), input.revision.clone())),
);
let preview = hydrate_catalog(base.as_ref().clone(), &catalogs, pins).map_err(|error| {
startup_failure(
StartupStage::Content,
"verify restored release revisions",
error,
)
})?;
for digest in self
.installed_content
.iter()
.chain(self.retained.iter().map(|item| &item.content_digest))
{
if !catalogs.contains_key(digest) {
return Err(startup_failure(
StartupStage::Content,
"verify restored pinned candidates",
RetainedCatalogError::CandidateUnavailable,
));
}
}
self.verify_activations(&catalogs, &preview)?;
rebuild_public_snapshot(
&self.startup,
&catalogs,
&preview,
embedded_manifest(),
self.tip_recipient.as_ref(),
)?;
Ok(())
}
fn verify_activations(
&self,
catalogs: &BTreeMap<ContentTreeDigest, Arc<ContentCatalog>>,
preview: &Arc<ContentCatalog>,
) -> Result<(), ProcessError> {
for activation in self.startup.activating.iter().cloned() {
rebuild_activating_publication(
activation,
&self.startup,
catalogs,
preview,
embedded_manifest(),
self.tip_recipient.as_ref(),
)
.map_err(|error| {
startup_failure(StartupStage::Content, "verify restored activation", error)
})?;
}
Ok(())
}
}
#[derive(Debug, thiserror::Error)]
enum StartupRecoveryError {
#[error("the activating publication has no durable site head")]
MissingSite,
#[error("the activating publication's retained content candidate is unavailable")]
MissingCandidate,
#[error("the activating publication's retained revisions could not be hydrated")]
Revisions(#[from] CatalogRetentionError),
#[error("the activating publication could not be reconstructed")]
Activation(#[from] PublicationActivationError),
}
fn rebuild_activating_publication(
activation: RecoverablePublicationActivation,
startup: &StartupSnapshotState,
catalogs: &BTreeMap<ContentTreeDigest, Arc<ContentCatalog>>,
preview: &ContentCatalog,
frontend: &'static FrontendAssetManifest,
tip_recipient: Option<&TipRecipientProjection>,
) -> Result<PreparedPublicationRecovery, StartupRecoveryError> {
let site = startup
.site
.clone()
.ok_or(StartupRecoveryError::MissingSite)?;
let mut catalog = catalogs
.get(&activation.content_digest)
.ok_or(StartupRecoveryError::MissingCandidate)?
.as_ref()
.clone();
catalog.retain_revisions_from(preview, startup.ledger.revision_keys())?;
Ok(PreparedPublicationRecovery::prepare(
activation,
Arc::new(catalog),
frontend,
&startup.ledger,
tip_recipient,
site,
)?)
}
fn compile_retained_catalogs(
store: &ContentCandidateStore,
compiler: &ContentCompiler,
) -> Result<BTreeMap<ContentTreeDigest, Arc<ContentCatalog>>, RetainedCatalogError> {
store
.load_all()
.map_err(RetainedCatalogError::Load)?
.into_iter()
.map(|candidate| {
let digest = candidate.digest;
let content = prepare_content(&candidate.tree).map_err(|source| {
RetainedCatalogError::Prepare {
digest: digest.clone(),
source,
}
})?;
match compiler.compile(&content) {
Ok(catalog) => Ok((digest, Arc::new(catalog))),
Err(source) => Err(RetainedCatalogError::Compile { digest, source }),
}
})
.collect()
}
fn hydrate_catalog(
mut base: ContentCatalog,
retained: &BTreeMap<ContentTreeDigest, Arc<ContentCatalog>>,
pins: impl IntoIterator<Item = (PostId, PostRevisionDigest)>,
) -> Result<Arc<ContentCatalog>, RetainedCatalogError> {
for (post_id, revision) in pins {
if base.get(&post_id, &revision).is_some() {
continue;
}
let source = retained
.values()
.find(|catalog| catalog.get(&post_id, &revision).is_some())
.ok_or_else(|| RetainedCatalogError::RevisionUnavailable {
post_id: post_id.clone(),
revision: revision.clone(),
})?;
base.retain_revisions_from(source, std::iter::once((post_id.clone(), revision.clone())))
.map_err(RetainedCatalogError::Retain)?;
}
Ok(Arc::new(base))
}
fn find_retained_public_snapshot(
retained: &BTreeMap<ContentTreeDigest, Arc<ContentCatalog>>,
ledger: &PublicLedgerProjection,
expected: &SiteSnapshotDigest,
frontend: &'static FrontendAssetManifest,
tip_recipient: Option<&TipRecipientProjection>,
) -> Result<SiteSnapshot, RetainedCatalogError> {
let pins = ledger
.published_posts()
.map(|published| (published.post_id.clone(), published.revision.clone()))
.collect::<Vec<_>>();
for base in retained.values() {
let catalog = hydrate_catalog(base.as_ref().clone(), retained, pins.clone())?;
let Ok(shell) = render_site_shell(catalog, frontend, ledger) else {
continue;
};
let shell = shell.bind_tip_recipient(tip_recipient.cloned());
let Ok(snapshot) = shell.into_snapshot() else {
continue;
};
if &snapshot.digest == expected {
return Ok(snapshot);
}
}
Err(RetainedCatalogError::PublicSnapshotUnavailable {
expected: expected.clone(),
})
}
fn rebuild_public_snapshot(
startup: &StartupSnapshotState,
retained: &BTreeMap<ContentTreeDigest, Arc<ContentCatalog>>,
preview_catalog: &Arc<ContentCatalog>,
frontend: &'static FrontendAssetManifest,
tip_recipient: Option<&TipRecipientProjection>,
) -> Result<SiteSnapshot, ProcessError> {
let Some(expected) = startup.site.as_ref() else {
let shell = render_site_shell(Arc::clone(preview_catalog), frontend, &startup.ledger)
.map_err(|error| {
startup_failure(StartupStage::Content, "render the site shell", error)
})?
.bind_tip_recipient(tip_recipient.cloned());
return shell.into_snapshot().map_err(|error| {
startup_failure(StartupStage::Content, "build the site snapshot", error)
});
};
find_retained_public_snapshot(
retained,
&startup.ledger,
&expected.digest,
frontend,
tip_recipient,
)
.map_err(|error| {
startup_failure(
StartupStage::Content,
"rebuild the approved public site snapshot",
error,
)
})
}
#[derive(Debug, thiserror::Error)]
enum RetainedCatalogError {
#[error("a required retained content candidate is unavailable")]
CandidateUnavailable,
#[error("retained content candidates could not be loaded")]
Load(#[source] ContentCandidateStoreError),
#[error("retained content candidate {digest} could not be prepared")]
Prepare {
digest: ContentTreeDigest,
#[source]
source: PrepareContentError,
},
#[error("retained content candidate {digest} could not be compiled")]
Compile {
digest: ContentTreeDigest,
#[source]
source: CatalogBuildError,
},
#[error("retained revision {revision} for post {post_id} is unavailable")]
RevisionUnavailable {
post_id: PostId,
revision: PostRevisionDigest,
},
#[error("a retained revision could not be installed in the recovery catalog")]
Retain(#[source] CatalogRetentionError),
#[error("no retained content candidate rebuilds durable site {expected}")]
PublicSnapshotUnavailable { expected: SiteSnapshotDigest },
}
struct ServingStateInput<'resources> {
database: &'resources DatabaseStore,
mail: &'resources MailConfiguration,
compiled: CompiledStartupContent,
frontend: &'static FrontendAssetManifest,
public_bind: std::net::SocketAddr,
admin_bind: AdminBind,
metrics_bind: std::net::SocketAddr,
backup: BackupHealth,
security: AdminSecurityState,
cancellation: CancellationToken,
candidate_store: &'resources ContentCandidateStore,
content_compiler: &'resources ContentCompiler,
source: SourceSyncHandle,
}
async fn prepare_serving_state(input: ServingStateInput<'_>) -> Result<ServingState, ProcessError> {
let ServingStateInput {
database,
mail,
compiled,
frontend,
public_bind,
admin_bind,
metrics_bind,
backup,
security,
cancellation,
candidate_store,
content_compiler,
source,
} = input;
database
.mail
.quarantine_interrupted(OffsetDateTime::now_utc())
.await
.map_err(|error| {
startup_failure(
StartupStage::Database,
"quarantine interrupted mail campaigns",
error,
)
})?;
let PreparedMail {
access: mail_access,
public_routes: mail_routes,
dispatcher: mail_dispatcher,
feedback: mail_feedback,
} = prepare_mail(
mail,
database,
compiled.catalog.publication.site.base_url.clone(),
)
.await
.map_err(|error| startup_failure(StartupStage::Configuration, "prepare mail runtime", error))?;
let tip_recipient = database
.profiles
.effective_tip_recipient()
.await
.map_err(|error| {
startup_failure(
StartupStage::Database,
"load the active tip recipient profile",
error,
)
})?;
let mut startup_state = database
.publications
.startup_snapshot_state()
.await
.map_err(|error| {
startup_failure(
StartupStage::Database,
"load the startup publication ledger",
error,
)
})?;
let ledger = startup_state.ledger.clone();
let retained_catalogs =
compile_retained_catalogs(candidate_store, content_compiler).map_err(|error| {
startup_failure(
StartupStage::Content,
"compile retained content candidates",
error,
)
})?;
let preview_pins = ledger
.published_posts()
.map(|published| (published.post_id.clone(), published.revision.clone()))
.chain(startup_state.scheduled.iter().map(|scheduled| {
let view = scheduled.publication.view();
(view.stable_post_id.clone(), view.pinned_post_digest.clone())
}))
.chain(startup_state.activating.iter().map(|activation| {
let view = activation.publication.view();
(view.stable_post_id.clone(), view.pinned_post_digest.clone())
}))
.collect::<Vec<_>>();
let preview_catalog = hydrate_catalog(
compiled.catalog.as_ref().clone(),
&retained_catalogs,
preview_pins,
)
.map_err(|error| {
startup_failure(
StartupStage::Content,
"hydrate the private preview catalog",
error,
)
})?;
let recovery = startup_state
.activating
.pop()
.map(|activation| {
rebuild_activating_publication(
activation,
&startup_state,
&retained_catalogs,
&preview_catalog,
frontend,
tip_recipient.as_ref(),
)
})
.transpose()
.map_err(|error| {
startup_failure(
StartupStage::Content,
"rebuild the activating publication candidate",
error,
)
})?;
let snapshot = rebuild_public_snapshot(
&startup_state,
&retained_catalogs,
&preview_catalog,
frontend,
tip_recipient.as_ref(),
)?;
let installed_site = database
.publications
.install_startup_snapshot(InstallStartupSnapshot {
expected: startup_state.site.clone(),
candidate_digest: snapshot.digest.clone(),
activated_at: OffsetDateTime::now_utc(),
source_commit: compiled.source_commit.clone(),
posts: compiled.observed_posts,
})
.await
.map_err(|error| {
startup_failure(
StartupStage::Database,
"install the startup site snapshot",
error,
)
})?;
let readiness = Readiness::default();
let (snapshots, activator) = snapshot_store(snapshot);
let mut publication_coordinator = PublicationCoordinator {
catalog: preview_catalog,
content_digest: compiled.content_digest,
candidates: Arc::new(retained_catalogs),
ledger,
site: installed_site,
activator,
store: database.publications.clone(),
profiles: database.profiles.clone(),
tip_recipient,
frontend,
source_commit: compiled.source_commit,
scheduled: startup_state
.scheduled
.drain(..)
.map(|scheduled| (scheduled.publication_id, scheduled))
.collect(),
scheduler_wakeup: Arc::new(tokio::sync::Notify::new()),
readiness: readiness.clone(),
cancellation: cancellation.clone(),
};
if let Some(recovery) = recovery {
publication_coordinator
.recover(recovery)
.await
.map_err(|error| {
startup_failure(
StartupStage::Database,
"recover the activating publication",
error,
)
})?;
}
let (publication_coordinator, publication_actor) =
publication_coordinator.into_actor(PUBLICATION_COORDINATOR_QUEUE_CAPACITY);
let protected_admin_router = runtime_admin_router(
publication_coordinator.clone(),
security,
database.profiles.clone(),
source,
MailUiState {
campaigns: database.mail.clone(),
subscribers: database.subscribers.clone(),
publications: publication_coordinator.clone(),
snapshots: snapshots.clone(),
access: mail_access,
},
);
let public_state = PublicState {
snapshots,
readiness: readiness.clone(),
};
let public_server = match mail_routes {
Some(routes) => {
PublicServer::bind_router(public_bind, public_router_with_routes(public_state, routes))
.await
}
None => PublicServer::bind(public_bind, public_state).await,
}
.map_err(|error| startup_failure(StartupStage::Listeners, "bind the public listener", error))?;
tracing::info!(bind = %public_server.local_addr, "public listener bound");
let admin_server = AdminServer::bind(admin_bind, protected_admin_router)
.await
.map_err(|error| {
startup_failure(StartupStage::Listeners, "bind the admin listener", error)
})?;
tracing::info!(
bind = %admin_server.local_addr,
"authenticated admin backend listener bound"
);
let metrics = Metrics::new(&database.health.metrics).map_err(|error| {
startup_failure(
StartupStage::Listeners,
"construct the application metrics registry",
error,
)
})?;
let metrics_server = MetricsServer::bind(metrics_bind, metrics.clone())
.await
.map_err(|error| {
startup_failure(StartupStage::Listeners, "bind the metrics listener", error)
})?;
tracing::info!(bind = %metrics_server.local_addr, "loopback metrics listener bound");
let metrics_collector = MetricsCollector::new(
metrics,
database.health.clone(),
backup,
tokio::runtime::Handle::current(),
);
Ok(ServingState {
mail_dispatcher,
mail_feedback,
metrics_server,
metrics_collector,
readiness,
publication_coordinator,
publication_actor,
public_server,
admin_server,
})
}
fn startup_failure(
stage: StartupStage,
operation: &'static str,
error: impl StdError + Send + Sync + 'static,
) -> ProcessError {
tracing::error!(%stage, operation, error = %error, "startup operation failed");
ApplicationError::Startup {
stage,
operation,
source: Box::new(error),
}
.into()
}
async fn start_database(host: &HostConfiguration) -> Result<StartedDatabase, ProcessError> {
let host = host.view();
restore::verify_startup_candidate(host.database.path, host.state_root)
.await
.map_err(|error| {
startup_failure(
StartupStage::Database,
"verify restored startup candidate",
error,
)
})?;
let database = match database::bootstrap(host.database).await {
Ok(database) => database,
Err(database::DatabaseStartupError::AlreadyOwned) => {
return Err(ProcessError::AlreadyRunning);
}
Err(error) => {
return Err(startup_failure(
StartupStage::Database,
"bootstrap the database",
error,
));
}
};
let (store, writer) = database.into_store(host.database.writer_queue_capacity.get());
let shutdown = CancellationToken::new();
let task = spawn_critical_task(CriticalTask::new(
CriticalTaskName::DatabaseWriter,
writer.run(shutdown.clone()),
));
if let Err(error) = identity_bootstrap::initialize_startup_identity(
&store,
host.identity_startup_bootstrap,
std::io::stdout(),
)
.await
{
let error = startup_failure(
StartupStage::Identity,
"apply the owner identity startup policy",
error,
);
return Err(close_writer_after_startup_failure(store, shutdown, task, error).await);
}
let providers = ConfiguredLoginProviders::new(true, true)
.expect("password and Nostr form a valid login-provider set");
let security = match AdminSecurityState::new(
host.admin_origin.clone(),
store.auth.clone(),
providers,
AdminSessionPolicy::default(),
Argon2idPolicy::v1(),
)
.await
{
Ok(security) => security,
Err(error) => {
let error = startup_failure(
StartupStage::Identity,
"initialize the admin authentication boundary",
error,
);
return Err(close_writer_after_startup_failure(store, shutdown, task, error).await);
}
};
Ok(StartedDatabase {
store,
shutdown,
task,
security,
})
}
async fn close_started_database(
database: StartedDatabase,
process_error: ProcessError,
) -> ProcessError {
close_writer_after_startup_failure(
database.store,
database.shutdown,
database.task,
process_error,
)
.await
}
async fn close_writer_after_startup_failure(
database: DatabaseStore,
shutdown: CancellationToken,
writer: JoinHandle<(CriticalTaskName, CriticalTaskCompletion)>,
process_error: ProcessError,
) -> ProcessError {
shutdown.cancel();
drop(database);
if let Some(error) = drained_task_failure(writer.await) {
tracing::error!(error = %error, "database writer failed during startup cleanup");
}
process_error
}
impl ApplicationRuntime {
fn with_parts(
readiness: Readiness,
cancellation: CancellationToken,
shutdown: ShutdownFuture,
critical_tasks: Vec<CriticalTask>,
) -> Self {
let mut supervisor = JoinSet::new();
for task in critical_tasks {
let span = task_span(task.name);
supervisor.spawn(task.instrument(span));
}
Self {
readiness,
cancellation,
database_shutdown: CancellationToken::new(),
shutdown,
critical_tasks: supervisor,
database_writer: None,
}
}
fn with_database_writer(
readiness: Readiness,
cancellation: CancellationToken,
database_shutdown: CancellationToken,
shutdown: ShutdownFuture,
critical_tasks: Vec<CriticalTask>,
database_writer: JoinHandle<(CriticalTaskName, CriticalTaskCompletion)>,
) -> Self {
let mut runtime = Self::with_parts(readiness, cancellation, shutdown, critical_tasks);
runtime.database_shutdown = database_shutdown;
runtime.database_writer = Some(database_writer);
runtime
}
async fn run_until_stop(mut self) -> Result<(), ApplicationError> {
self.readiness.mark_ready();
let failure = if self.critical_tasks.is_empty() && self.database_writer.is_none() {
self.shutdown.await.err()
} else {
tokio::select! {
biased;
shutdown = &mut self.shutdown => shutdown.err(),
completion = self.critical_tasks.join_next(), if !self.critical_tasks.is_empty() => {
Some(unexpected_task_failure(completion))
}
completion = wait_for_database_writer(&mut self.database_writer) => {
self.database_writer.take();
Some(unexpected_task_failure(Some(completion)))
}
}
};
self.readiness.mark_not_ready();
self.cancellation.cancel();
let drain_failure = drain_critical_tasks(&mut self.critical_tasks).await;
self.database_shutdown.cancel();
let database_failure = drain_database_writer(&mut self.database_writer).await;
match first_shutdown_failure([failure, drain_failure, database_failure]) {
Some(error) => Err(error),
None => Ok(()),
}
}
}
fn spawn_critical_task(
task: CriticalTask,
) -> JoinHandle<(CriticalTaskName, CriticalTaskCompletion)> {
let span = task_span(task.name);
tokio::spawn(task.instrument(span))
}
fn first_shutdown_failure(failures: [Option<ApplicationError>; 3]) -> Option<ApplicationError> {
let mut failures = failures.into_iter().flatten();
let first = failures.next();
for failure in failures {
tracing::error!(error = %failure, "additional failure during shutdown");
}
first
}
async fn wait_for_database_writer(
writer: &mut Option<JoinHandle<(CriticalTaskName, CriticalTaskCompletion)>>,
) -> Result<(CriticalTaskName, CriticalTaskCompletion), JoinError> {
match writer.as_mut() {
Some(writer) => writer.await,
None => std::future::pending().await,
}
}
async fn drain_database_writer(
writer: &mut Option<JoinHandle<(CriticalTaskName, CriticalTaskCompletion)>>,
) -> Option<ApplicationError> {
let writer = writer.take()?;
drained_task_failure(writer.await)
}
async fn drain_critical_tasks(
critical_tasks: &mut JoinSet<(CriticalTaskName, CriticalTaskCompletion)>,
) -> Option<ApplicationError> {
let mut first_failure = None;
while let Some(completion) = critical_tasks.join_next().await {
if let Some(failure) = drained_task_failure(completion) {
if first_failure.is_none() {
first_failure = Some(failure);
} else {
tracing::error!(error = %failure, "additional critical task failure during shutdown");
}
}
}
first_failure
}
fn unexpected_task_failure(
completion: Option<Result<(CriticalTaskName, CriticalTaskCompletion), JoinError>>,
) -> ApplicationError {
match completion {
Some(Ok((task, completion))) => task_completion_failure(task, completion)
.unwrap_or(ApplicationError::CriticalTaskExited { task }),
Some(Err(source)) => ApplicationError::TaskSupervisor { source },
None => ApplicationError::TaskSupervisorEmpty,
}
}
fn drained_task_failure(
completion: Result<(CriticalTaskName, CriticalTaskCompletion), JoinError>,
) -> Option<ApplicationError> {
match completion {
Ok((task, completion)) => task_completion_failure(task, completion),
Err(source) => Some(ApplicationError::TaskSupervisor { source }),
}
}
fn task_completion_failure(
task: CriticalTaskName,
completion: CriticalTaskCompletion,
) -> Option<ApplicationError> {
match completion {
CriticalTaskCompletion::Returned(Ok(())) => None,
CriticalTaskCompletion::Returned(Err(source)) => {
Some(ApplicationError::CriticalTaskFailed { task, source })
}
CriticalTaskCompletion::Panicked(message) => {
Some(ApplicationError::CriticalTaskPanicked { task, message })
}
}
}
struct CriticalTask {
name: CriticalTaskName,
future: CriticalTaskFuture,
}
impl CriticalTask {
fn new<Future, Error>(name: CriticalTaskName, future: Future) -> Self
where
Future: std::future::Future<Output = Result<(), Error>> + Send + 'static,
Error: Into<CriticalTaskFailure>,
{
Self {
name,
future: Box::pin(async move { future.await.map_err(Into::into) }),
}
}
}
enum CriticalTaskCompletion {
Returned(CriticalTaskResult),
Panicked(Box<str>),
}
impl Future for CriticalTask {
type Output = (CriticalTaskName, CriticalTaskCompletion);
fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Self::Output> {
let task = self.get_mut();
match catch_unwind(AssertUnwindSafe(|| task.future.as_mut().poll(context))) {
Ok(Poll::Pending) => Poll::Pending,
Ok(Poll::Ready(result)) => {
Poll::Ready((task.name, CriticalTaskCompletion::Returned(result)))
}
Err(payload) => Poll::Ready((
task.name,
CriticalTaskCompletion::Panicked(panic_message(payload.as_ref())),
)),
}
}
}
fn panic_message(payload: &(dyn std::any::Any + Send)) -> Box<str> {
if let Some(message) = payload.downcast_ref::<String>() {
return message.clone().into_boxed_str();
}
if let Some(message) = payload.downcast_ref::<&'static str>() {
return Box::from(*message);
}
Box::from("non-string panic payload")
}
#[cfg(unix)]
fn install_termination_signal() -> Result<ShutdownFuture, ApplicationError> {
use tokio::signal::unix::{SignalKind, signal};
let interrupt =
signal(SignalKind::interrupt()).map_err(|source| ApplicationError::SignalRegistration {
signal: ShutdownSignal::Interrupt,
source,
})?;
let terminate =
signal(SignalKind::terminate()).map_err(|source| ApplicationError::SignalRegistration {
signal: ShutdownSignal::Terminate,
source,
})?;
Ok(Box::pin(wait_for_unix_termination(interrupt, terminate)))
}
#[cfg(unix)]
async fn wait_for_unix_termination(
mut interrupt: tokio::signal::unix::Signal,
mut terminate: tokio::signal::unix::Signal,
) -> Result<(), ApplicationError> {
let closed_signal = tokio::select! {
received = interrupt.recv() => {
if received.is_some() {
return Ok(());
}
ShutdownSignal::Interrupt
}
received = terminate.recv() => {
if received.is_some() {
return Ok(());
}
ShutdownSignal::Terminate
}
};
Err(ApplicationError::SignalStreamClosed {
signal: closed_signal,
})
}
#[cfg(windows)]
fn install_termination_signal() -> Result<ShutdownFuture, ApplicationError> {
let interrupt = tokio::signal::windows::ctrl_c().map_err(|source| {
ApplicationError::SignalRegistration {
signal: ShutdownSignal::Interrupt,
source,
}
})?;
Ok(Box::pin(wait_for_windows_termination(interrupt)))
}
#[cfg(windows)]
async fn wait_for_windows_termination(
mut interrupt: tokio::signal::windows::CtrlC,
) -> Result<(), ApplicationError> {
if interrupt.recv().await.is_some() {
Ok(())
} else {
Err(ApplicationError::SignalStreamClosed {
signal: ShutdownSignal::Interrupt,
})
}
}
#[cfg(not(any(unix, windows)))]
fn install_termination_signal() -> Result<ShutdownFuture, ApplicationError> {
Err(ApplicationError::SignalPlatformUnsupported)
}
#[cfg(test)]
mod tests {
use std::{
cell::Cell,
fs,
path::PathBuf,
sync::{
Arc,
atomic::{AtomicBool, AtomicUsize, Ordering},
},
};
use axum::http::{Method, StatusCode};
use k256::schnorr::SigningKey;
use markdown_compiler::DefaultPostTipPolicy;
use serde::{Serialize, de::DeserializeOwned};
use tokio::{
io::{AsyncReadExt as _, AsyncWriteExt as _},
net::{TcpSocket, TcpStream},
sync::oneshot,
};
use maincopy_shared::auth::{
AdminAuditEventId, AdminScope, AgentCredentialId, InstanceId, UserId,
};
use crate::{
admin::test_support::{ADMIN_AUTHORITY, ADMIN_ORIGIN, agent_authorization},
cli::BootstrapCredential,
config::{ConfigurationValidationCode, IdentityStartupBootstrap},
domain::{
auth::{
NostrPublicKey,
store::{
AdminMutationKey, AuditPrincipalReference, BootstrapIdentity,
MutationAuditContext, NewHumanCredential, RegisterAgentCredential,
},
},
mail::runtime::{MailReviewAccessError, MailStartupError},
},
render::render_bound_post_preview,
};
use super::*;
const VALID_PUBLICATION: &str = "[site]\n\
title = \"Pinned startup source\"\n\
base_url = \"https://startup.example.test\"\n\
description = \"Startup configuration test.\"\n\
[author]\n\
name = \"Startup Tester\"\n";
const DURABLE_POST_ID: &str = "11111111-1111-4111-8111-111111111111";
const DURABLE_PUBLICATION_ID: &str = "aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa";
const LIVE_RELOAD_TEST_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(15);
const TEST_OWNER_NOSTR_PUBLIC_KEY: &str =
"63fe6318dc58583cfe16810f86dd09e18bfd76aabc24a0081ce2856f330504ed";
const TEST_AGENT_SECRET: [u8; 32] = [3_u8; 32];
const DURABLE_POST: &str = "+++\n\
id = \"11111111-1111-4111-8111-111111111111\"\n\
title = \"Durable publication\"\n\
slug = \"durable-publication\"\n\
authored_at = 2026-08-29T12:00:00Z\n\
description = \"A publication restored from SQLite.\"\n\
+++\n\
# Durable publication\n\n\
Durable article body.\n";
fn reserve_loopback_port() -> TcpSocket {
let reservation = TcpSocket::new_v4().unwrap();
reservation.set_reuseaddr(true).unwrap();
reservation.bind("127.0.0.1:0".parse().unwrap()).unwrap();
reservation
}
fn startup_host_source(extra: &str, public_bind: &str) -> String {
startup_host_source_with_listeners(extra, public_bind, "127.0.0.1:0", "127.0.0.1:0")
}
fn startup_host_source_with_listeners(
extra: &str,
public_bind: &str,
admin_bind: &str,
metrics_bind: &str,
) -> String {
format!(
"[paths]\n\
content_root = \"content\"\n\
state_root = \"state\"\n\
runtime_root = \"run\"\n\
[public]\n\
bind = \"{public_bind}\"\n\
{extra}\n\
[metrics]\n\
bind = \"{metrics_bind}\"\n\
[admin]\n\
bind = \"{admin_bind}\"\n\
origin = \"https://admin.example.test\"\n"
)
}
fn startup_fixture(
host_source: &str,
publication_source: &str,
) -> (tempfile::TempDir, PathBuf, PathBuf) {
let root = tempfile::tempdir().unwrap();
let content_root = root.path().join("content");
fs::create_dir(&content_root).unwrap();
let publication_path = content_root.join("publication.toml");
fs::write(&publication_path, publication_source).unwrap();
let config_path = root.path().join("maincopy.toml");
fs::write(
&config_path,
startup_host_source(host_source, "127.0.0.1:0"),
)
.unwrap();
(root, config_path, publication_path)
}
fn write_durable_post(content_root: &Path) {
fs::create_dir(content_root.join("posts")).unwrap();
fs::write(
content_root.join("posts/durable-publication.md"),
DURABLE_POST,
)
.unwrap();
}
struct AdminTestClient {
address: std::net::SocketAddr,
client: reqwest::Client,
signing_key: SigningKey,
}
fn admin_client(application: &Application) -> AdminTestClient {
AdminTestClient {
address: application.admin_addr,
client: reqwest::Client::builder()
.no_proxy()
.timeout(std::time::Duration::from_secs(5))
.build()
.unwrap(),
signing_key: SigningKey::from_bytes(&TEST_AGENT_SECRET).unwrap(),
}
}
async fn admin_get(admin: &AdminTestClient, path: &str) -> reqwest::Response {
let idempotency_key = uuid::Uuid::new_v4().to_string();
let authorization = agent_authorization(
&admin.signing_key,
&Method::GET,
path,
&[],
Some(&idempotency_key),
);
admin
.client
.get(format!("http://{}{path}", admin.address))
.header(reqwest::header::HOST, ADMIN_AUTHORITY)
.header("origin", ADMIN_ORIGIN)
.header(reqwest::header::AUTHORIZATION, authorization)
.header(
maincopy_shared::publication::IDEMPOTENCY_KEY_HEADER,
idempotency_key,
)
.send()
.await
.unwrap()
}
async fn admin_post_json(
admin: &AdminTestClient,
path: &str,
idempotency_key: &str,
value: &impl Serialize,
) -> reqwest::Response {
admin_post_body(
admin,
path,
idempotency_key,
None,
serde_json::to_vec(value).unwrap(),
)
.await
}
async fn admin_post_body(
admin: &AdminTestClient,
path: &str,
idempotency_key: &str,
request_id: Option<&str>,
body: impl AsRef<[u8]>,
) -> reqwest::Response {
let body = body.as_ref().to_vec();
let authorization = agent_authorization(
&admin.signing_key,
&Method::POST,
path,
&body,
Some(idempotency_key),
);
let mut request = admin
.client
.post(format!("http://{}{path}", admin.address))
.header(reqwest::header::HOST, ADMIN_AUTHORITY)
.header("origin", ADMIN_ORIGIN)
.header(reqwest::header::AUTHORIZATION, authorization)
.header(
maincopy_shared::publication::IDEMPOTENCY_KEY_HEADER,
idempotency_key,
)
.header(reqwest::header::CONTENT_TYPE, "application/json");
if let Some(request_id) = request_id {
request = request.header("x-request-id", request_id);
}
request.body(body).send().await.unwrap()
}
async fn admin_json<ResponseType: DeserializeOwned>(
response: reqwest::Response,
) -> ResponseType {
let bytes = response.bytes().await.unwrap();
serde_json::from_slice(&bytes).unwrap()
}
async fn admin_text(response: reqwest::Response) -> String {
response.text().await.unwrap()
}
async fn admin_preview_digest(
admin: &AdminTestClient,
post_id: uuid::Uuid,
) -> maincopy_shared::publication::PreviewDigest {
use maincopy_shared::{
posts::POSTS_PATH,
publication::{PREVIEW_DIGEST_HEADER, PreviewDigest},
};
let response = admin_get(admin, &format!("{POSTS_PATH}/{post_id}/preview")).await;
assert_eq!(response.status(), StatusCode::OK);
let encoded = response
.headers()
.get(PREVIEW_DIGEST_HEADER)
.unwrap()
.to_str()
.unwrap();
PreviewDigest::parse(encoded).unwrap()
}
async fn stop_built_application(application: Application) {
let Application {
_process_lock: process_lock,
_database: database,
publication_coordinator,
mut runtime,
public_addr: _,
admin_addr: _,
} = application;
runtime.cancellation.cancel();
assert!(
drain_critical_tasks(&mut runtime.critical_tasks)
.await
.is_none()
);
runtime.database_shutdown.cancel();
assert!(
drain_database_writer(&mut runtime.database_writer)
.await
.is_none()
);
drop(publication_coordinator);
drop(database);
drop(process_lock);
tokio::task::yield_now().await;
}
async fn bootstrap_test_identity(startup: &StartupConfiguration) {
let host = startup._host.view();
let database = database::bootstrap(host.database).await.unwrap();
let (store, writer) = database.into_store(host.database.writer_queue_capacity.get());
let shutdown = CancellationToken::new();
let writer_shutdown = shutdown.clone();
let writer = tokio::spawn(async move {
writer.run(writer_shutdown).await.unwrap();
});
if store
.auth
.identity_state()
.await
.unwrap()
.bootstrap_required
{
let owner_user_id = UserId::from_uuid(uuid::Uuid::new_v4());
store
.auth
.bootstrap_identity(BootstrapIdentity {
instance_id: InstanceId::from_uuid(uuid::Uuid::new_v4()),
owner_user_id,
credential: NewHumanCredential::Nostr {
public_key: NostrPublicKey::parse(TEST_OWNER_NOSTR_PUBLIC_KEY).unwrap(),
},
configured_providers: ConfiguredLoginProviders::new(true, true).unwrap(),
occurred_at: OffsetDateTime::now_utc(),
audit_event_id: AdminAuditEventId::from_uuid(uuid::Uuid::new_v4()),
})
.await
.unwrap();
let signing_key = SigningKey::from_bytes(&TEST_AGENT_SECRET).unwrap();
store
.auth
.register_agent_credential(RegisterAgentCredential {
credential_id: AgentCredentialId::from_uuid(uuid::Uuid::new_v4()),
owner_user_id,
issuer_user_id: owner_user_id,
public_key: NostrPublicKey::from_bytes(
signing_key.verifying_key().to_bytes().into(),
)
.unwrap(),
label: "startup integration agent".into(),
scopes: AdminScope::PUBLISHER.into_iter().collect(),
created_at: OffsetDateTime::now_utc(),
expires_at: None,
audit: MutationAuditContext {
audit_event_id: AdminAuditEventId::from_uuid(uuid::Uuid::new_v4()),
principal: AuditPrincipalReference::Offline {
user_id: Some(owner_user_id),
},
request_id: None,
idempotency_key: AdminMutationKey(uuid::Uuid::new_v4()),
},
})
.await
.unwrap();
}
drop(store);
shutdown.cancel();
writer.await.unwrap();
}
async fn build_test_application(
startup: StartupConfiguration,
) -> Result<Application, ProcessError> {
bootstrap_test_identity(&startup).await;
Application::build(startup).await
}
#[test]
#[cfg(target_os = "linux")]
fn content_discovery_runs_once_receives_effective_limits_and_validation_reuses_owned_bytes() {
let host_source = "[content]\n\
publication_file_bytes = 512\n\
post_file_bytes = 1536\n\
asset_file_bytes = 2560\n\
total_tree_bytes = 3584\n\
entries = 90\n\
depth = 7\n\
path_bytes = 400\n";
let (_root, arguments, publication_path) = startup_fixture(host_source, VALID_PUBLICATION);
let discovery_calls = Cell::new(0);
let observed_limits = Cell::new(None);
let startup =
StartupConfiguration::load_with_discovery(arguments, |content_root, limits| {
discovery_calls.set(discovery_calls.get() + 1);
observed_limits.set(Some(limits));
let tree = discover_content_tree(content_root, limits)?;
fs::write(&publication_path, "not the discovered publication").unwrap();
Ok(tree)
})
.unwrap();
assert_eq!(discovery_calls.get(), 1);
let limits = observed_limits.get().unwrap();
assert_eq!(limits.publication_file_bytes.get(), 512);
assert_eq!(limits.post_file_bytes.get(), 1536);
assert_eq!(limits.asset_file_bytes.get(), 2560);
assert_eq!(limits.total_tree_bytes.get(), 3584);
assert_eq!(limits.entries.get(), 90);
assert_eq!(limits.depth.get(), 7);
assert_eq!(limits.path_bytes.get(), 400);
assert!(
startup
._content_tree
.publication
.source
.contains("Pinned startup source")
);
assert_eq!(
startup.content.view().publication.site.title.as_str(),
"Pinned startup source"
);
}
#[test]
fn host_configuration_failure_prevents_content_discovery() {
let root = tempfile::tempdir().unwrap();
let arguments = root.path().join("missing-maincopy.toml");
let discovery_calls = Cell::new(0);
let result = StartupConfiguration::load_with_discovery(arguments, |_, _| {
discovery_calls.set(discovery_calls.get() + 1);
unreachable!("content discovery must follow host configuration")
});
assert!(matches!(result, Err(ProcessError::Configuration(_))));
assert_eq!(discovery_calls.get(), 0);
}
#[test]
#[cfg(target_os = "linux")]
fn content_validation_failure_has_stable_exit() {
let (_root, arguments, _) = startup_fixture("", "unknown = true\n");
let discovery_calls = Cell::new(0);
let result = StartupConfiguration::load_with_discovery(arguments, |root, limits| {
discovery_calls.set(discovery_calls.get() + 1);
discover_content_tree(root, limits)
});
let Err(error) = result else {
panic!("invalid content must fail startup");
};
assert!(matches!(error, ProcessError::Validation(_)));
assert_eq!(error.exit(), ProcessExit::Validation);
assert_eq!(discovery_calls.get(), 1);
}
#[test]
#[cfg(target_os = "linux")]
fn authored_tips_do_not_require_a_payment_provider() {
let publication = format!(
"{VALID_PUBLICATION}[tips]\n\
enabled = true\n"
);
let (_root, arguments, _) = startup_fixture("", &publication);
let startup =
StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
assert_eq!(
startup.content.view().publication.tips,
DefaultPostTipPolicy::Enabled
);
}
#[test]
#[cfg(target_os = "linux")]
fn removed_provider_configuration_is_rejected_without_opening_credentials() {
let host = "[lightning]\n\
provider = \"lexe\"\n\
network = \"signet\"\n\
credentials = { source = \"file\", path = \"must-not-open.json\" }\n";
let (root, arguments, _) = startup_fixture(host, VALID_PUBLICATION);
let credential_path = root.path().join("must-not-open.json");
let Err(error) =
StartupConfiguration::load_with_discovery(arguments, discover_content_tree)
else {
panic!("removed provider configuration must fail host validation");
};
let ProcessError::Configuration(errors) = error else {
panic!("removed provider configuration must fail host validation");
};
assert_eq!(
errors.diagnostics()[0].code,
ConfigurationValidationCode::HostTomlInvalid
);
assert!(!credential_path.exists());
}
#[tokio::test]
async fn production_identity_policy_refuses_credential_output_and_accepts_offline_nostr_bootstrap()
{
let (_root, config_path, _) = startup_fixture(
"[identity]\nstartup_bootstrap = \"require_existing\"\n",
VALID_PUBLICATION,
);
let host = HostConfigurationLoader::from_process_working_directory()
.unwrap()
.load(&config_path)
.unwrap();
let database = database::bootstrap(host.view().database).await.unwrap();
let (store, writer) = database.into_store(host.view().database.writer_queue_capacity.get());
let cancellation = CancellationToken::new();
let running = tokio::spawn(writer.run(cancellation.clone()));
let mut output = Vec::new();
let failure = identity_bootstrap::initialize_startup_identity(
&store,
IdentityStartupBootstrap::RequireExisting,
&mut output,
)
.await
.unwrap_err();
assert!(matches!(
failure,
identity_bootstrap::GeneratedOwnerBootstrapError::ExistingIdentityRequired
));
assert!(output.is_empty());
assert!(
store
.auth
.identity_state()
.await
.unwrap()
.bootstrap_required
);
cancellation.cancel();
running.await.unwrap().unwrap();
drop(store);
let startup =
StartupConfiguration::load_with_discovery(config_path.clone(), discover_content_tree)
.unwrap();
let failure = match Application::build(startup).await {
Ok(application) => {
stop_built_application(application).await;
panic!("production startup must refuse an uninitialized identity");
}
Err(error) => error,
};
assert!(matches!(
failure,
ProcessError::Application(ApplicationError::Startup {
stage: StartupStage::Identity,
operation: "apply the owner identity startup policy",
..
})
));
let public_key = NostrPublicKey::parse(
"f9308a019258c31049344f85f89d5229b531c845836f99b08601f113bce036f9",
)
.unwrap();
identity_bootstrap::bootstrap_owner(
config_path.clone(),
BootstrapCredential::Nostr { public_key },
)
.await
.unwrap();
let startup =
StartupConfiguration::load_with_discovery(config_path, discover_content_tree).unwrap();
let application = Application::build(startup).await.unwrap();
let owner = application
._database
.auth
.users_page(None, 2)
.await
.unwrap();
assert_eq!(owner.items.len(), 1);
assert!(owner.items[0].has_nostr);
assert!(!owner.items[0].has_password);
let mut output = Vec::new();
assert!(
!identity_bootstrap::initialize_startup_identity(
&application._database,
IdentityStartupBootstrap::RequireExisting,
&mut output
)
.await
.unwrap()
);
assert!(output.is_empty());
stop_built_application(application).await;
}
#[tokio::test]
#[cfg(target_os = "linux")]
async fn unbootstrapped_identity_generates_owner_and_starts_the_application() {
let (root, _, _) = startup_fixture("", VALID_PUBLICATION);
let config_path = root.path().join("maincopy.toml");
let reservation = reserve_loopback_port();
let public_addr = reservation.local_addr().unwrap();
fs::write(
&config_path,
startup_host_source("", &public_addr.to_string()),
)
.unwrap();
let arguments = config_path;
let startup =
StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
let application = Application::build(startup).await.unwrap();
let identity = application._database.auth.identity_state().await.unwrap();
assert!(!identity.bootstrap_required);
assert!(identity.instance.is_some());
let users = application
._database
.auth
.users_page(None, 2)
.await
.unwrap();
assert_eq!(users.items.len(), 1);
let owner = &users.items[0];
assert!(owner.has_password);
assert!(!owner.has_nostr);
assert!(
owner
.roles
.contains(&maincopy_shared::auth::UserRole::Owner)
);
let credentials = application
._database
.auth
.user_credentials(owner.user_id)
.await
.unwrap()
.unwrap();
assert!(matches!(
credentials.as_slice(),
[crate::domain::auth::store::StoredHumanCredential::Password {
username,
..
}] if username.as_str() == "owner"
));
assert!(
!identity_bootstrap::initialize_startup_identity(
&application._database,
IdentityStartupBootstrap::GenerateOwner,
std::io::sink()
)
.await
.unwrap()
);
stop_built_application(application).await;
drop(reservation);
}
#[tokio::test]
#[cfg(target_os = "linux")]
async fn application_build_serves_public_site_and_protected_admin_backend() {
let (root, arguments, _) = startup_fixture("", VALID_PUBLICATION);
let startup =
StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
let application = build_test_application(startup).await.unwrap();
assert!(matches!(
StartupHostConfiguration::load(root.path().join("maincopy.toml")),
Err(ProcessError::AlreadyRunning)
));
assert_eq!(application.runtime.critical_tasks.len(), 8);
assert!(application.runtime.database_writer.is_some());
let client = reqwest::Client::builder()
.no_proxy()
.timeout(std::time::Duration::from_secs(2))
.build()
.unwrap();
let response = client
.get(format!("http://{}/", application.public_addr))
.send()
.await
.unwrap();
assert_eq!(response.status(), reqwest::StatusCode::OK);
let html = response.text().await.unwrap();
assert!(html.contains("Pinned startup source"));
assert!(html.contains("No posts have been published yet."));
let protected = client
.get(format!(
"http://{}/api/admin/v1/capabilities",
application.admin_addr
))
.header(reqwest::header::HOST, "admin.example.test")
.send()
.await
.unwrap();
assert_eq!(protected.status(), reqwest::StatusCode::UNAUTHORIZED);
assert_eq!(
protected.headers()[reqwest::header::CACHE_CONTROL],
"private, no-store"
);
stop_built_application(application).await;
assert!(StartupHostConfiguration::load(root.path().join("maincopy.toml")).is_ok());
}
#[tokio::test]
#[cfg(target_os = "linux")]
async fn admin_publication_route_activates_the_public_site_and_replays_success() {
use maincopy_shared::publication::{
PUBLICATIONS_PATH, PublishNowRequest, PublishNowResponse,
};
let (root, arguments, _) = startup_fixture("", VALID_PUBLICATION);
let content_root = root.path().join("content");
write_durable_post(&content_root);
let startup =
StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
let application = build_test_application(startup).await.unwrap();
let admin = admin_client(&application);
let post_id = uuid::Uuid::parse_str(DURABLE_POST_ID).unwrap();
let request = PublishNowRequest {
post_id,
preview_digest: admin_preview_digest(&admin, post_id).await,
expected_revision: None,
scheduled_for: None,
};
let creation_key = "bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb";
let response = admin_post_json(&admin, PUBLICATIONS_PATH, creation_key, &request).await;
assert_eq!(response.status(), StatusCode::OK);
let published: PublishNowResponse = admin_json(response).await;
assert_eq!(published.post_id, request.post_id);
assert!(published.revision.starts_with("post-b3-v1-"));
assert!(published.site_digest.starts_with("site-b3-v1-"));
assert_eq!(published.site_version, 2);
let replay = admin_post_body(
&admin,
PUBLICATIONS_PATH,
creation_key,
None,
serde_json::to_string_pretty(&request).unwrap(),
)
.await;
assert_eq!(replay.status(), StatusCode::OK);
assert_eq!(admin_json::<PublishNowResponse>(replay).await, published);
let request_id = "67e55044-10b1-426f-9247-bb680e5fe0c8";
let malformed_key = uuid::Uuid::new_v4().to_string();
let malformed = admin_post_body(
&admin,
PUBLICATIONS_PATH,
&malformed_key,
Some(request_id),
"{",
)
.await;
assert_eq!(malformed.status(), StatusCode::BAD_REQUEST);
let error: serde_json::Value = admin_json(malformed).await;
assert_eq!(error["error"]["code"], "invalid_request_body");
assert_eq!(error["error"]["request_id"], request_id);
let oversized_key = uuid::Uuid::new_v4().to_string();
let oversized = admin_post_body(
&admin,
PUBLICATIONS_PATH,
&oversized_key,
Some(request_id),
" ".repeat(4 * 1024 + 1),
)
.await;
assert_eq!(oversized.status(), StatusCode::PAYLOAD_TOO_LARGE);
assert_eq!(oversized.headers()["x-request-id"], request_id);
let error: serde_json::Value = admin_json(oversized).await;
assert_eq!(error["error"]["code"], "request_body_too_large");
assert_eq!(error["error"]["request_id"], request_id);
assert_eq!(
admin_get(&admin, "/api/admin/v1/capabilities")
.await
.status(),
StatusCode::OK
);
let public = reqwest::get(format!(
"http://{}/posts/durable-publication",
application.public_addr
))
.await
.unwrap();
assert_eq!(public.status(), reqwest::StatusCode::OK);
assert!(
public
.text()
.await
.unwrap()
.contains("Durable article body.")
);
stop_built_application(application).await;
}
#[tokio::test]
#[cfg(target_os = "linux")]
async fn edited_markdown_updates_private_preview_until_explicit_publication_approval() {
use maincopy_shared::{
posts::{ListPostsResponse, POSTS_PATH, PostPublicationState},
publication::{PUBLICATIONS_PATH, PublishNowRequest, PublishNowResponse},
};
let (root, arguments, _) = startup_fixture("", VALID_PUBLICATION);
let content_root = root.path().join("content");
write_durable_post(&content_root);
let post_path = content_root.join("posts/durable-publication.md");
let startup =
StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
let application = build_test_application(startup).await.unwrap();
let admin = admin_client(&application);
let response = admin_post_json(
&admin,
PUBLICATIONS_PATH,
"bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb",
&PublishNowRequest {
post_id: uuid::Uuid::parse_str(DURABLE_POST_ID).unwrap(),
preview_digest: admin_preview_digest(
&admin,
uuid::Uuid::parse_str(DURABLE_POST_ID).unwrap(),
)
.await,
expected_revision: None,
scheduled_for: None,
},
)
.await;
assert_eq!(response.status(), StatusCode::OK);
let published: PublishNowResponse = admin_json(response).await;
let public_url = format!(
"http://{}/posts/durable-publication",
application.public_addr
);
let public = reqwest::Client::builder().no_proxy().build().unwrap();
let initial = public.get(&public_url).send().await.unwrap();
let initial_etag = initial.headers()[reqwest::header::ETAG].clone();
assert!(
initial
.text()
.await
.unwrap()
.contains("Durable article body.")
);
let edited = DURABLE_POST.replace(
"Durable article body.",
"This edit appeared without restarting Maincopy.",
);
fs::write(&post_path, &edited).unwrap();
let edited_summary = tokio::time::timeout(LIVE_RELOAD_TEST_TIMEOUT, async {
loop {
let posts: ListPostsResponse =
admin_json(admin_get(&admin, POSTS_PATH).await).await;
if let Some(summary) = posts.posts.into_iter().find(|summary| {
summary.post_id.to_string() == DURABLE_POST_ID
&& summary.publication_state == PostPublicationState::UnpublishedChange
&& summary.revision != published.revision
}) {
break summary;
}
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
})
.await
.expect("a stable Markdown edit must enter the private preview catalog");
let preview_url = format!("{POSTS_PATH}/{}/preview", edited_summary.post_id);
let preview = admin_get(&admin, &preview_url).await;
assert_eq!(preview.status(), StatusCode::OK);
let preview_body = admin_text(preview).await;
assert!(preview_body.contains("This edit appeared without restarting Maincopy."));
assert!(!preview_body.contains("Durable article body."));
let still_pinned = public.get(&public_url).send().await.unwrap();
assert_eq!(still_pinned.headers()[reqwest::header::ETAG], initial_etag);
let still_pinned_body = still_pinned.text().await.unwrap();
assert!(still_pinned_body.contains("Durable article body."));
assert!(!still_pinned_body.contains("This edit appeared without restarting Maincopy."));
fs::write(&post_path, "+++\ninvalid = true\n").unwrap();
tokio::time::sleep(std::time::Duration::from_millis(1_500)).await;
let posts: ListPostsResponse = admin_json(admin_get(&admin, POSTS_PATH).await).await;
let after_invalid = posts
.posts
.into_iter()
.find(|summary| summary.post_id == edited_summary.post_id)
.unwrap();
assert_eq!(after_invalid.revision, edited_summary.revision);
assert_eq!(
after_invalid.publication_state,
PostPublicationState::UnpublishedChange
);
let preview = admin_get(&admin, &preview_url).await;
assert_eq!(preview.status(), StatusCode::OK);
assert!(
admin_text(preview)
.await
.contains("This edit appeared without restarting Maincopy.")
);
let after_invalid_public = public.get(&public_url).send().await.unwrap();
assert_eq!(
after_invalid_public.headers()[reqwest::header::ETAG],
initial_etag
);
let after_invalid_public_body = after_invalid_public.text().await.unwrap();
assert!(after_invalid_public_body.contains("Durable article body."));
assert!(
!after_invalid_public_body.contains("This edit appeared without restarting Maincopy.")
);
fs::write(&post_path, &edited).unwrap();
let response = admin_post_json(
&admin,
PUBLICATIONS_PATH,
"cccccccc-cccc-4ccc-8ccc-cccccccccccc",
&PublishNowRequest {
post_id: edited_summary.post_id,
preview_digest: admin_preview_digest(&admin, edited_summary.post_id).await,
expected_revision: Some(edited_summary.revision.clone()),
scheduled_for: None,
},
)
.await;
assert_eq!(response.status(), StatusCode::OK);
let approved: PublishNowResponse = admin_json(response).await;
assert_eq!(approved.revision, edited_summary.revision);
let updated_public = public.get(&public_url).send().await.unwrap();
assert_ne!(
updated_public.headers()[reqwest::header::ETAG],
initial_etag
);
let updated_body = updated_public.text().await.unwrap();
assert!(updated_body.contains("This edit appeared without restarting Maincopy."));
assert!(!updated_body.contains("Durable article body."));
stop_built_application(application).await;
}
#[tokio::test]
#[cfg(target_os = "linux")]
async fn a_new_markdown_file_becomes_publishable_without_restarting() {
use maincopy_shared::{
posts::{ListPostsResponse, POSTS_PATH},
publication::{PUBLICATIONS_PATH, PublishNowRequest},
};
let (root, arguments, _) = startup_fixture("", VALID_PUBLICATION);
let content_root = root.path().join("content");
let startup =
StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
let application = build_test_application(startup).await.unwrap();
let admin = admin_client(&application);
write_durable_post(&content_root);
let summary = tokio::time::timeout(LIVE_RELOAD_TEST_TIMEOUT, async {
loop {
let posts: ListPostsResponse =
admin_json(admin_get(&admin, POSTS_PATH).await).await;
if let Some(post) = posts
.posts
.into_iter()
.find(|post| post.post_id.to_string() == DURABLE_POST_ID)
{
break post;
}
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
})
.await
.expect("a stable new Markdown file must enter the live catalog");
let response = admin_post_json(
&admin,
PUBLICATIONS_PATH,
"bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb",
&PublishNowRequest {
post_id: summary.post_id,
preview_digest: admin_preview_digest(&admin, summary.post_id).await,
expected_revision: Some(summary.revision),
scheduled_for: None,
},
)
.await;
assert_eq!(response.status(), StatusCode::OK);
let body = reqwest::get(format!(
"http://{}/posts/durable-publication",
application.public_addr
))
.await
.unwrap()
.text()
.await
.unwrap();
assert!(body.contains("Durable article body."));
stop_built_application(application).await;
}
#[tokio::test]
#[cfg(target_os = "linux")]
async fn restart_keeps_public_revision_and_aliases_pinned_until_stopped_edit_is_approved() {
use maincopy_shared::{
posts::{ListPostsResponse, POSTS_PATH, PostPublicationState},
publication::{PUBLICATIONS_PATH, PublishNowRequest, PublishNowResponse},
};
let (root, arguments, _) = startup_fixture("", VALID_PUBLICATION);
let content_root = root.path().join("content");
write_durable_post(&content_root);
let durable_post_path = content_root.join("posts/durable-publication.md");
fs::write(
&durable_post_path,
DURABLE_POST.replace(
"slug = \"durable-publication\"\n",
"slug = \"durable-publication\"\naliases = [\"original-durable-alias\"]\n",
),
)
.unwrap();
let startup =
StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
let application = build_test_application(startup).await.unwrap();
let admin = admin_client(&application);
let response = admin_post_json(
&admin,
PUBLICATIONS_PATH,
"bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb",
&PublishNowRequest {
post_id: uuid::Uuid::parse_str(DURABLE_POST_ID).unwrap(),
preview_digest: admin_preview_digest(
&admin,
uuid::Uuid::parse_str(DURABLE_POST_ID).unwrap(),
)
.await,
expected_revision: None,
scheduled_for: None,
},
)
.await;
assert_eq!(response.status(), StatusCode::OK);
let published: PublishNowResponse = admin_json(response).await;
let public = reqwest::Client::builder()
.no_proxy()
.redirect(reqwest::redirect::Policy::none())
.build()
.unwrap();
let initial_public_url = format!(
"http://{}/posts/durable-publication",
application.public_addr
);
let initial_alias_url = format!(
"http://{}/posts/original-durable-alias",
application.public_addr
);
let initial_alias = public.get(&initial_alias_url).send().await.unwrap();
assert_eq!(initial_alias.status(), StatusCode::PERMANENT_REDIRECT);
assert_eq!(
initial_alias.headers()[reqwest::header::LOCATION],
"https://startup.example.test/posts/durable-publication"
);
let initial = public.get(&initial_public_url).send().await.unwrap();
let initial_etag = initial.headers()[reqwest::header::ETAG].clone();
assert!(
initial
.text()
.await
.unwrap()
.contains("Durable article body.")
);
stop_built_application(application).await;
let stopped_edit = DURABLE_POST
.replace(
"slug = \"durable-publication\"\n",
"slug = \"durable-publication\"\naliases = [\"edited-durable-alias\"]\n",
)
.replace(
"Durable article body.",
"This edit was made while Maincopy was stopped.",
);
fs::write(&durable_post_path, stopped_edit).unwrap();
let arguments = root.path().join("maincopy.toml");
let startup =
StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
let restarted = build_test_application(startup).await.unwrap();
let public_url = format!("http://{}/posts/durable-publication", restarted.public_addr);
let pinned = public.get(&public_url).send().await.unwrap();
assert_eq!(pinned.headers()[reqwest::header::ETAG], initial_etag);
let pinned_body = pinned.text().await.unwrap();
assert!(pinned_body.contains("Durable article body."));
assert!(!pinned_body.contains("This edit was made while Maincopy was stopped."));
let pinned_alias_url = format!(
"http://{}/posts/original-durable-alias",
restarted.public_addr
);
let pinned_alias = public.get(&pinned_alias_url).send().await.unwrap();
assert_eq!(pinned_alias.status(), StatusCode::PERMANENT_REDIRECT);
assert_eq!(
pinned_alias.headers()[reqwest::header::LOCATION],
"https://startup.example.test/posts/durable-publication"
);
let edited_alias_url = format!(
"http://{}/posts/edited-durable-alias",
restarted.public_addr
);
assert_eq!(
public.get(&edited_alias_url).send().await.unwrap().status(),
StatusCode::NOT_FOUND
);
let admin = admin_client(&restarted);
let posts: ListPostsResponse = admin_json(admin_get(&admin, POSTS_PATH).await).await;
let edited_summary = posts
.posts
.into_iter()
.find(|summary| summary.post_id.to_string() == DURABLE_POST_ID)
.unwrap();
assert_eq!(
edited_summary.publication_state,
PostPublicationState::UnpublishedChange
);
assert_ne!(edited_summary.revision, published.revision);
let preview = admin_get(
&admin,
&format!("{POSTS_PATH}/{}/preview", edited_summary.post_id),
)
.await;
assert_eq!(preview.status(), StatusCode::OK);
let preview_body = admin_text(preview).await;
assert!(preview_body.contains("This edit was made while Maincopy was stopped."));
assert!(!preview_body.contains("Durable article body."));
let response = admin_post_json(
&admin,
PUBLICATIONS_PATH,
"cccccccc-cccc-4ccc-8ccc-cccccccccccc",
&PublishNowRequest {
post_id: edited_summary.post_id,
preview_digest: admin_preview_digest(&admin, edited_summary.post_id).await,
expected_revision: Some(edited_summary.revision.clone()),
scheduled_for: None,
},
)
.await;
assert_eq!(response.status(), StatusCode::OK);
let approved: PublishNowResponse = admin_json(response).await;
assert_eq!(approved.revision, edited_summary.revision);
let updated = public.get(&public_url).send().await.unwrap();
assert_ne!(updated.headers()[reqwest::header::ETAG], initial_etag);
let updated_body = updated.text().await.unwrap();
assert!(updated_body.contains("This edit was made while Maincopy was stopped."));
assert!(!updated_body.contains("Durable article body."));
assert_eq!(
public.get(&pinned_alias_url).send().await.unwrap().status(),
StatusCode::NOT_FOUND
);
let updated_alias = public.get(&edited_alias_url).send().await.unwrap();
assert_eq!(updated_alias.status(), StatusCode::PERMANENT_REDIRECT);
assert_eq!(
updated_alias.headers()[reqwest::header::LOCATION],
"https://startup.example.test/posts/durable-publication"
);
stop_built_application(restarted).await;
}
#[tokio::test]
#[cfg(target_os = "linux")]
async fn restart_serves_the_durable_published_revision() {
use maincopy_shared::{
posts::{ListPostsResponse, POSTS_PATH, PostPublicationState},
publication::{PUBLICATIONS_PATH, PublishNowRequest, PublishNowResponse},
};
let (root, arguments, _) = startup_fixture("", VALID_PUBLICATION);
let content_root = root.path().join("content");
write_durable_post(&content_root);
let startup =
StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
let application = build_test_application(startup).await.unwrap();
let admin = admin_client(&application);
let response = admin_post_json(
&admin,
PUBLICATIONS_PATH,
"bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb",
&PublishNowRequest {
post_id: uuid::Uuid::parse_str(DURABLE_POST_ID).unwrap(),
preview_digest: admin_preview_digest(
&admin,
uuid::Uuid::parse_str(DURABLE_POST_ID).unwrap(),
)
.await,
expected_revision: None,
scheduled_for: None,
},
)
.await;
assert_eq!(response.status(), StatusCode::OK);
let published: PublishNowResponse = admin_json(response).await;
let published_at = published.published_at.unwrap();
let public = reqwest::Client::builder().no_proxy().build().unwrap();
let initial_url = format!(
"http://{}/posts/durable-publication",
application.public_addr
);
let initial = public.get(&initial_url).send().await.unwrap();
let initial_etag = initial.headers()[reqwest::header::ETAG].clone();
let initial_html = initial.text().await.unwrap();
assert!(initial_html.contains("Durable article body."));
assert!(initial_html.contains(&published_at.to_string()));
stop_built_application(application).await;
let arguments = root.path().join("maincopy.toml");
let startup =
StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
let application = build_test_application(startup).await.unwrap();
let restarted_admin = admin_client(&application);
let posts: ListPostsResponse =
admin_json(admin_get(&restarted_admin, POSTS_PATH).await).await;
let durable = posts
.posts
.into_iter()
.find(|summary| summary.post_id == published.post_id)
.unwrap();
assert_eq!(durable.publication_state, PostPublicationState::Published);
assert_eq!(durable.revision, published.revision);
let response = public
.get(format!(
"http://{}/posts/durable-publication",
application.public_addr
))
.send()
.await
.unwrap();
assert_eq!(response.status(), reqwest::StatusCode::OK);
assert_eq!(response.headers()[reqwest::header::ETAG], initial_etag);
let html = response.text().await.unwrap();
assert!(html.contains("Durable article body."));
assert!(html.contains(&published_at.to_string()));
stop_built_application(application).await;
}
#[tokio::test]
#[cfg(target_os = "linux")]
async fn offline_restore_requires_exact_scheduled_and_blocked_candidate_inputs() {
use sqlx::{ConnectOptions as _, Connection as _};
let (root, arguments, _) = startup_fixture("", VALID_PUBLICATION);
let content_root = root.path().join("content");
write_durable_post(&content_root);
let startup =
StartupConfiguration::load_with_discovery(arguments.clone(), discover_content_tree)
.unwrap();
let database_path = startup._host.view().database.path.to_owned();
let limits = startup._host.view().content_limits;
let state_root = root.path().join("state");
stop_built_application(build_test_application(startup).await.unwrap()).await;
let approved = discover_content_tree(&content_root, limits).unwrap();
let approved_digest = approved.digest();
let catalog =
Arc::new(compile_content_catalog(&prepare_content(&approved).unwrap()).unwrap());
let post_id = PostId::parse(DURABLE_POST_ID).unwrap();
let revision = catalog.current_post(&post_id).unwrap().revision.clone();
let preview = render_bound_post_preview(
&catalog,
embedded_manifest(),
&post_id,
None,
"/api/admin/v1/preview-assets/retained-fixture",
None,
)
.unwrap()
.unwrap();
let candidates = ContentCandidateStore::open(&state_root, limits).unwrap();
fs::write(
content_root.join("publication.toml"),
VALID_PUBLICATION.replace("Pinned startup source", "Another publication title"),
)
.unwrap();
let same_revision = discover_content_tree(&content_root, limits).unwrap();
assert_eq!(
compile_content_catalog(&prepare_content(&same_revision).unwrap())
.unwrap()
.current_post(&post_id)
.unwrap()
.revision,
revision,
);
let alternate_digest = candidates.retain(&same_revision).unwrap();
assert_ne!(alternate_digest, approved_digest);
fs::write(
content_root.join("posts/durable-publication.md"),
DURABLE_POST.replace("Durable article body.", "A later unapproved revision."),
)
.unwrap();
let later = discover_content_tree(&content_root, limits).unwrap();
candidates.retain(&later).unwrap();
let mut connection = sqlx::sqlite::SqliteConnectOptions::new()
.filename(&database_path)
.foreign_keys(true)
.connect()
.await
.unwrap();
sqlx::query(
"INSERT INTO canonical_publications (\
publication_id, creation_key, command_kind, stable_post_id, pinned_post_digest, \
state, version, scheduled_at_ns, content_tree_digest, accepted_preview_digest\
) VALUES (?, ?, 'scheduled', ?, ?, 'scheduled', 1, ?, ?, ?)",
)
.bind(
uuid::Uuid::parse_str(DURABLE_PUBLICATION_ID)
.unwrap()
.as_bytes()
.as_slice(),
)
.bind(uuid::Uuid::new_v4().as_bytes().as_slice())
.bind(post_id.as_uuid().as_bytes().as_slice())
.bind(revision.as_bytes().as_slice())
.bind(1_900_000_000_000_000_000_i64)
.bind(approved_digest.as_bytes().as_slice())
.bind(preview.digest.as_bytes().as_slice())
.execute(&mut connection)
.await
.unwrap();
connection.close().await.unwrap();
for (state, version, activation, reason) in [
("scheduled", 1_i64, None, None),
(
"blocked",
3_i64,
Some(1_900_000_000_000_000_000_i64),
Some("preview_changed"),
),
] {
let mut connection = sqlx::sqlite::SqliteConnectOptions::new()
.filename(&database_path)
.connect()
.await
.unwrap();
sqlx::query(
"UPDATE canonical_publications SET state = ?, version = ?, \
activation_at_ns = ?, block_reason = ?",
)
.bind(state)
.bind(version)
.bind(activation)
.bind(reason)
.execute(&mut connection)
.await
.unwrap();
connection.close().await.unwrap();
verify_offline_content_without_mutation(&database_path, &state_root, limits)
.await
.unwrap();
fs::remove_file(
state_root.join(format!("content-candidates/{approved_digest}.candidate")),
)
.unwrap();
let error =
verify_offline_content_without_mutation(&database_path, &state_root, limits)
.await
.unwrap_err();
assert!(matches!(
error,
ProcessError::Application(ApplicationError::Startup {
operation: "verify restored pinned candidates",
..
})
));
fs::remove_file(
state_root.join(format!("content-candidates/{alternate_digest}.candidate")),
)
.unwrap();
let error =
verify_offline_content_without_mutation(&database_path, &state_root, limits)
.await
.unwrap_err();
assert!(matches!(
error,
ProcessError::Application(ApplicationError::Startup {
operation: "verify restored release revisions",
..
})
));
candidates.retain(&approved).unwrap();
candidates.retain(&same_revision).unwrap();
}
}
async fn verify_offline_content_without_mutation(
database_path: &Path,
state_root: &Path,
limits: ContentTreeLimits,
) -> Result<(), ProcessError> {
let before = blake3::hash(&fs::read(database_path).unwrap());
let inspected = database::restore::inspect(database_path).await.unwrap();
let result = verify_restore_content(&inspected.store, state_root, limits).await;
inspected.close().await;
assert_eq!(blake3::hash(&fs::read(database_path).unwrap()), before);
result
}
#[tokio::test]
#[cfg(target_os = "linux")]
async fn startup_preflights_exact_activation_before_indexing_new_content_or_binding() {
use sqlx::{ConnectOptions as _, Connection as _};
let (root, arguments, _) = startup_fixture("", VALID_PUBLICATION);
let content_root = root.path().join("content");
write_durable_post(&content_root);
let startup =
StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
let database_path = startup._host.view().database.path.to_owned();
let limits = startup._host.view().content_limits;
let state_root = root.path().join("state");
stop_built_application(build_test_application(startup).await.unwrap()).await;
let arguments = root.path().join("maincopy.toml");
let startup =
StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
let content_digest = startup._content_tree.digest();
let catalog = Arc::new(compile_content_catalog(&startup.content).unwrap());
let post_id = PostId::parse(DURABLE_POST_ID).unwrap();
let rendered = catalog.current_post(&post_id).unwrap();
let revision = rendered.revision.clone();
let accepted_preview_digest = render_bound_post_preview(
&catalog,
embedded_manifest(),
&post_id,
None,
"/api/admin/v1/preview-assets/recovery-fixture",
None,
)
.unwrap()
.unwrap()
.digest;
let activation_at = OffsetDateTime::from_unix_timestamp(1_777_734_400).unwrap();
let candidate_ledger = PublicLedgerProjection::empty()
.with_published(PublishedPostRevision::new(
post_id,
revision.clone(),
activation_at,
))
.unwrap();
let shell = render_site_shell(Arc::clone(&catalog), embedded_manifest(), &candidate_ledger)
.unwrap();
let candidate = shell.into_snapshot().unwrap();
let candidate_digest = candidate.digest.clone();
let activation_at_ns = i64::try_from(activation_at.unix_timestamp_nanos()).unwrap();
let publication_id = uuid::Uuid::parse_str(DURABLE_PUBLICATION_ID)
.unwrap()
.into_bytes();
let creation_key = uuid::Uuid::parse_str("bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb")
.unwrap()
.into_bytes();
let stable_post_id = uuid::Uuid::parse_str(DURABLE_POST_ID).unwrap().into_bytes();
let mut connection = sqlx::sqlite::SqliteConnectOptions::new()
.filename(&database_path)
.foreign_keys(true)
.connect()
.await
.unwrap();
sqlx::query(
"INSERT INTO canonical_publications (\
publication_id, creation_key, command_kind, stable_post_id, pinned_post_digest, state, version, \
scheduled_at_ns, activation_at_ns, activation_site_digest, content_tree_digest, \
accepted_preview_digest\
) VALUES (?, ?, 'immediate', ?, ?, 'activating', 2, ?, ?, ?, ?, ?)",
)
.bind(publication_id.as_slice())
.bind(creation_key.as_slice())
.bind(stable_post_id.as_slice())
.bind(revision.as_bytes().as_slice())
.bind(activation_at_ns)
.bind(activation_at_ns)
.bind(candidate_digest.as_bytes().as_slice())
.bind(content_digest.as_bytes().as_slice())
.bind(accepted_preview_digest.as_bytes().as_slice())
.execute(&mut connection)
.await
.unwrap();
let original_site: Vec<u8> =
sqlx::query_scalar("SELECT current_site_digest FROM site_state WHERE singleton = 1")
.fetch_one(&mut connection)
.await
.unwrap();
sqlx::query("UPDATE canonical_publications SET activation_site_digest = ?")
.bind([0xab_u8; 32].as_slice())
.execute(&mut connection)
.await
.unwrap();
connection.close().await.unwrap();
drop(startup);
let error = verify_offline_content_without_mutation(&database_path, &state_root, limits)
.await
.unwrap_err();
assert!(matches!(
error,
ProcessError::Application(ApplicationError::Startup {
stage: StartupStage::Content,
operation: "verify restored activation",
..
})
));
fs::write(
content_root.join("posts/durable-publication.md"),
DURABLE_POST.replace("Durable article body.", "Later private revision."),
)
.unwrap();
let startup = StartupConfiguration::load_with_discovery(
root.path().join("maincopy.toml"),
discover_content_tree,
)
.unwrap();
let error = match build_test_application(startup).await {
Ok(_) => panic!("an unreproducible activation must reject startup"),
Err(error) => error,
};
assert!(matches!(
error,
ProcessError::Application(ApplicationError::Startup {
stage: StartupStage::Content,
operation: "rebuild the activating publication candidate",
..
})
));
let mut connection = sqlx::sqlite::SqliteConnectOptions::new()
.filename(&database_path)
.connect()
.await
.unwrap();
let revisions: Vec<Vec<u8>> =
sqlx::query_scalar("SELECT revision_digest FROM post_revisions")
.fetch_all(&mut connection)
.await
.unwrap();
assert_eq!(revisions, vec![revision.as_bytes().to_vec()]);
let durable_site: Vec<u8> =
sqlx::query_scalar("SELECT current_site_digest FROM site_state WHERE singleton = 1")
.fetch_one(&mut connection)
.await
.unwrap();
assert_eq!(durable_site, original_site);
let state: String = sqlx::query_scalar("SELECT state FROM canonical_publications")
.fetch_one(&mut connection)
.await
.unwrap();
assert_eq!(state, "activating");
sqlx::query("UPDATE canonical_publications SET activation_site_digest = ?")
.bind(candidate_digest.as_bytes().as_slice())
.execute(&mut connection)
.await
.unwrap();
connection.close().await.unwrap();
verify_offline_content_without_mutation(&database_path, &state_root, limits)
.await
.unwrap();
let startup = StartupConfiguration::load_with_discovery(
root.path().join("maincopy.toml"),
discover_content_tree,
)
.unwrap();
let application = build_test_application(startup).await.unwrap();
{
let projection = application.publication_coordinator.read();
assert_eq!(projection.site.digest, candidate_digest);
assert_eq!(projection.ledger.len(), 1);
let current = projection
.catalog
.current_post(&PostId::parse(DURABLE_POST_ID).unwrap())
.unwrap();
assert_ne!(current.revision, revision);
assert!(
current
.article
.identity_html
.contains("Later private revision.")
);
}
let response = reqwest::Client::builder()
.no_proxy()
.timeout(std::time::Duration::from_secs(2))
.build()
.unwrap()
.get(format!(
"http://{}/posts/durable-publication",
application.public_addr
))
.send()
.await
.unwrap();
assert_eq!(response.status(), reqwest::StatusCode::OK);
let html = response.text().await.unwrap();
assert!(html.contains("Durable article body."));
assert!(html.contains(&activation_at.to_string()));
stop_built_application(application).await;
let mut connection = sqlx::sqlite::SqliteConnectOptions::new()
.filename(&database_path)
.connect()
.await
.unwrap();
let (state, current_digest, published_at_ns): (String, Vec<u8>, i64) = sqlx::query_as(
"SELECT state, current_published_digest, published_at_ns \
FROM canonical_publications WHERE publication_id = ?",
)
.bind(publication_id.as_slice())
.fetch_one(&mut connection)
.await
.unwrap();
assert_eq!(state, "published");
assert_eq!(current_digest, revision.as_bytes());
assert_eq!(published_at_ns, activation_at_ns);
connection.close().await.unwrap();
}
#[tokio::test]
#[cfg(target_os = "linux")]
async fn mail_credential_failure_releases_database_process_and_listener_ownership() {
use std::{io::Write as _, os::unix::fs::OpenOptionsExt as _};
let (root, config_path, _) = startup_fixture("", VALID_PUBLICATION);
let reservations = [
reserve_loopback_port(),
reserve_loopback_port(),
reserve_loopback_port(),
];
let addresses = reservations
.each_ref()
.map(|socket| socket.local_addr().unwrap());
let host_source = startup_host_source_with_listeners(
"",
&addresses[0].to_string(),
&addresses[1].to_string(),
&addresses[2].to_string(),
);
fs::write(
&config_path,
format!(
"{host_source}\n\
[mail]\n\
mode = \"ses\"\n\
sender = \"newsletter@example.com\"\n\
region = \"us-east-1\"\n\
configuration_set = \"newsletter\"\n\
credential_file = \"private-mail-credential.json\"\n\
control_signing_key_file = \"unused-control.key\"\n"
),
)
.unwrap();
let credential_path = root.path().join("private-mail-credential.json");
let mut credential = fs::OpenOptions::new()
.write(true)
.create_new(true)
.mode(0o600)
.open(&credential_path)
.unwrap();
credential
.write_all(br#"{"access_key_id":"AKIDEXAMPLE","secret_access_key":"private-fixture-secret-material","unexpected":"invalid"}"#)
.unwrap();
drop(credential);
let startup =
StartupConfiguration::load_with_discovery(config_path.clone(), discover_content_tree)
.unwrap();
bootstrap_test_identity(&startup).await;
let error = match Application::build(startup).await {
Ok(_) => panic!("malformed SES credentials must fail application construction"),
Err(error) => error,
};
assert!(matches!(
&error,
ProcessError::Application(ApplicationError::Startup {
stage: StartupStage::Configuration,
operation: "prepare mail runtime",
source,
}) if matches!(source.downcast_ref::<MailStartupError>(),
Some(MailStartupError::Review(MailReviewAccessError::CredentialInvalid)))
));
let diagnostic = format!("{error}\n{error:?}");
for private in [
credential_path.to_str().unwrap(),
"private-mail-credential.json",
"private-fixture-secret-material",
] {
assert!(!diagnostic.contains(private));
}
let configuration = HostConfigurationLoader::from_process_working_directory()
.unwrap()
.load(&config_path)
.unwrap();
let host = configuration.view();
let process_lock = ProcessLock::acquire(host.runtime_root).unwrap();
let database = database::bootstrap(host.database).await.unwrap();
for address in addresses {
let listener = tokio::net::TcpListener::bind(address).await.unwrap();
drop(listener);
}
database.close().await.unwrap();
drop(process_lock);
drop(reservations);
}
#[tokio::test]
#[cfg(target_os = "linux")]
async fn listener_failure_releases_the_public_port_and_database_ownership() {
let (root, _, _) = startup_fixture("", VALID_PUBLICATION);
let config_path = root.path().join("maincopy.toml");
let reservation = reserve_loopback_port();
let public_addr = reservation.local_addr().unwrap();
let occupied = tokio::net::TcpListener::bind(public_addr).await.unwrap();
let public_bind = public_addr.to_string();
fs::write(&config_path, startup_host_source("", &public_bind)).unwrap();
let arguments = config_path.clone();
let startup =
StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
let error = match build_test_application(startup).await {
Ok(_) => panic!("an occupied public address must fail listener binding"),
Err(error) => error,
};
assert!(matches!(
error,
ProcessError::Application(ApplicationError::Startup {
stage: StartupStage::Listeners,
..
})
));
drop(occupied);
let arguments = config_path;
let startup =
StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
let application = build_test_application(startup).await.unwrap();
assert_eq!(application.public_addr, public_addr);
stop_built_application(application).await;
drop(reservation);
}
#[tokio::test]
#[cfg(target_os = "linux")]
async fn admin_listener_failure_releases_public_listener_and_database_ownership() {
let (root, _, _) = startup_fixture("", VALID_PUBLICATION);
let config_path = root.path().join("maincopy.toml");
let reservation = reserve_loopback_port();
let admin_addr = reservation.local_addr().unwrap();
let occupied = tokio::net::TcpListener::bind(admin_addr).await.unwrap();
fs::write(
&config_path,
startup_host_source_with_listeners(
"",
"127.0.0.1:0",
&admin_addr.to_string(),
"127.0.0.1:0",
),
)
.unwrap();
let arguments = config_path.clone();
let startup =
StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
let error = match build_test_application(startup).await {
Ok(_) => panic!("an occupied admin address must fail listener binding"),
Err(error) => error,
};
assert!(matches!(
error,
ProcessError::Application(ApplicationError::Startup {
stage: StartupStage::Listeners,
..
})
));
drop(occupied);
let arguments = config_path;
let startup =
StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
let application = build_test_application(startup).await.unwrap();
assert_eq!(application.admin_addr, admin_addr);
stop_built_application(application).await;
drop(reservation);
}
#[tokio::test]
#[cfg(target_os = "linux")]
async fn database_failure_prevents_application_build() {
use std::os::unix::fs::PermissionsExt as _;
use sqlx::{ConnectOptions as _, Connection as _};
let (_root, arguments, _) = startup_fixture("", VALID_PUBLICATION);
let startup =
StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
let database_path = startup._host.view().database.path.to_owned();
let database_parent = database_path.parent().unwrap();
fs::create_dir_all(database_parent).unwrap();
fs::set_permissions(database_parent, fs::Permissions::from_mode(0o700)).unwrap();
let mut foreign = sqlx::sqlite::SqliteConnectOptions::new()
.filename(&database_path)
.create_if_missing(true)
.connect()
.await
.unwrap();
sqlx::query("PRAGMA application_id = 7")
.execute(&mut foreign)
.await
.unwrap();
foreign.close().await.unwrap();
fs::set_permissions(&database_path, fs::Permissions::from_mode(0o600)).unwrap();
let error = match Application::build(startup).await {
Ok(_) => panic!("a foreign database must fail before listener binding"),
Err(error) => error,
};
assert!(matches!(
error,
ProcessError::Application(ApplicationError::Startup {
stage: StartupStage::Database,
..
})
));
}
async fn begin_public_request(address: std::net::SocketAddr, path: &str) -> TcpStream {
let mut stream = TcpStream::connect(address).await.unwrap();
let head = format!(
"GET {path} HTTP/1.1\r\nHost: localhost\r\nContent-Length: 1\r\nExpect: 100-continue\r\nConnection: close\r\n\r\n"
);
stream.write_all(head.as_bytes()).await.unwrap();
let mut interim = Vec::new();
tokio::time::timeout(std::time::Duration::from_secs(2), async {
while !interim.ends_with(b"\r\n\r\n") {
assert!(interim.len() < 1024);
interim.push(stream.read_u8().await.unwrap());
}
})
.await
.unwrap();
assert_eq!(
std::str::from_utf8(&interim).unwrap(),
"HTTP/1.1 100 Continue\r\n\r\n"
);
stream
}
async fn finish_public_request(mut stream: TcpStream) -> String {
stream.write_all(b"x").await.unwrap();
let mut response = String::new();
tokio::time::timeout(
std::time::Duration::from_secs(2),
stream.take(64 * 1024).read_to_string(&mut response),
)
.await
.unwrap()
.unwrap();
response
}
#[tokio::test]
async fn shutdown_finishes_an_accepted_public_request_before_closing_the_real_writer() {
let (_root, config, _) = startup_fixture("", VALID_PUBLICATION);
let reservation = reserve_loopback_port();
let address = reservation.local_addr().unwrap();
fs::write(&config, startup_host_source("", &address.to_string())).unwrap();
let startup =
StartupConfiguration::load_with_discovery(config, discover_content_tree).unwrap();
let mut application = build_test_application(startup).await.unwrap();
assert_eq!(application.public_addr, address);
let readiness = application.runtime.readiness.clone();
let cancellation = application.runtime.cancellation.clone();
let database_shutdown = application.runtime.database_shutdown.clone();
let (stop, stopped) = oneshot::channel();
application.runtime.shutdown = Box::pin(async {
stopped.await.unwrap();
Ok(())
});
let running = tokio::spawn(application.run_until_stop());
let request = begin_public_request(address, "/health/live").await;
stop.send(()).unwrap();
tokio::time::timeout(std::time::Duration::from_secs(2), cancellation.cancelled())
.await
.unwrap();
assert!(!readiness.is_ready());
assert!(!database_shutdown.is_cancelled());
assert!(!running.is_finished());
let response = finish_public_request(request).await;
assert!(response.starts_with("HTTP/1.1 200 OK"));
assert!(response.contains(r#"{"status":"live"}"#));
tokio::time::timeout(std::time::Duration::from_secs(2), running)
.await
.unwrap()
.unwrap()
.unwrap();
assert!(database_shutdown.is_cancelled());
let rebound = tokio::net::TcpListener::bind(address).await.unwrap();
assert_eq!(rebound.local_addr().unwrap(), address);
drop(rebound);
drop(reservation);
}
#[tokio::test]
async fn supervised_task_failure_makes_pending_readiness_fail_while_liveness_stays_live() {
let (_root, config, _) = startup_fixture("", VALID_PUBLICATION);
let startup =
StartupConfiguration::load_with_discovery(config, discover_content_tree).unwrap();
let mut application = build_test_application(startup).await.unwrap();
let address = application.public_addr;
let cancellation = application.runtime.cancellation.clone();
let (fail, failed) = oneshot::channel();
application.runtime.critical_tasks.spawn(CriticalTask::new(
CriticalTaskName::Worker,
async {
failed.await.unwrap();
Ok::<(), CriticalTaskFailure>(())
},
));
let running = tokio::spawn(application.run_until_stop());
let ready = begin_public_request(address, "/health/ready").await;
let live = begin_public_request(address, "/health/live").await;
fail.send(()).unwrap();
tokio::time::timeout(std::time::Duration::from_secs(2), cancellation.cancelled())
.await
.unwrap();
let ready = finish_public_request(ready).await;
let live = finish_public_request(live).await;
assert!(ready.starts_with("HTTP/1.1 503 Service Unavailable"));
assert!(ready.contains(r#"{"status":"not_ready"}"#));
assert!(live.starts_with("HTTP/1.1 200 OK"));
assert!(live.contains(r#"{"status":"live"}"#));
let failure = tokio::time::timeout(std::time::Duration::from_secs(2), running)
.await
.unwrap()
.unwrap()
.unwrap_err();
assert!(matches!(
failure,
ApplicationError::CriticalTaskExited {
task: CriticalTaskName::Worker
}
));
}
#[tokio::test]
async fn shutdown_signal_marks_unready_cancels_and_drains_every_task() {
let readiness = Readiness::default();
let cancellation = CancellationToken::new();
let drained = Arc::new(AtomicUsize::new(0));
let observed_order = Arc::new(AtomicUsize::new(0));
let (shutdown_tx, shutdown_rx) = oneshot::channel();
let tasks = (0..2)
.map(|_| {
let readiness = readiness.clone();
let cancellation = cancellation.clone();
let drained = Arc::clone(&drained);
let observed_order = Arc::clone(&observed_order);
CriticalTask::new(CriticalTaskName::Worker, async move {
cancellation.cancelled().await;
if !readiness.is_ready() {
observed_order.fetch_add(1, Ordering::SeqCst);
}
drained.fetch_add(1, Ordering::SeqCst);
Ok::<(), CriticalTaskFailure>(())
})
})
.collect();
let application = ApplicationRuntime::with_parts(
readiness.clone(),
cancellation.clone(),
Box::pin(async move {
let _ = shutdown_rx.await;
Ok(())
}),
tasks,
);
let running = tokio::spawn(application.run_until_stop());
while !readiness.is_ready() {
tokio::task::yield_now().await;
}
assert!(readiness.is_ready());
assert!(shutdown_tx.send(()).is_ok());
assert!(running.await.unwrap().is_ok());
assert!(!readiness.is_ready());
assert!(cancellation.is_cancelled());
assert_eq!(drained.load(Ordering::SeqCst), 2);
assert_eq!(observed_order.load(Ordering::SeqCst), 2);
}
#[tokio::test]
async fn shutdown_waits_for_the_database_writer_to_close() {
let readiness = Readiness::default();
let cancellation = CancellationToken::new();
let database_shutdown = CancellationToken::new();
let writer_shutdown = database_shutdown.clone();
let (shutdown_tx, shutdown_rx) = oneshot::channel();
let (writer_started_tx, writer_started_rx) = oneshot::channel();
let (writer_release_tx, writer_release_rx) = oneshot::channel();
let writer = CriticalTask::new(CriticalTaskName::DatabaseWriter, async move {
writer_shutdown.cancelled().await;
let _ = writer_started_tx.send(());
let _ = writer_release_rx.await;
Ok::<(), CriticalTaskFailure>(())
});
let application = ApplicationRuntime::with_database_writer(
readiness.clone(),
cancellation,
database_shutdown,
Box::pin(async move {
let _ = shutdown_rx.await;
Ok(())
}),
Vec::new(),
spawn_critical_task(writer),
);
let running = tokio::spawn(application.run_until_stop());
while !readiness.is_ready() {
tokio::task::yield_now().await;
}
assert!(shutdown_tx.send(()).is_ok());
assert!(writer_started_rx.await.is_ok());
tokio::task::yield_now().await;
assert!(!running.is_finished());
assert!(writer_release_tx.send(()).is_ok());
assert!(running.await.unwrap().is_ok());
assert!(!readiness.is_ready());
}
#[tokio::test]
async fn shutdown_drains_producers_before_stopping_the_database_writer() {
let readiness = Readiness::default();
let cancellation = CancellationToken::new();
let producer_shutdown = cancellation.clone();
let database_shutdown = CancellationToken::new();
let writer_shutdown = database_shutdown.clone();
let (shutdown_tx, shutdown_rx) = oneshot::channel();
let (producer_cancelled_tx, producer_cancelled_rx) = oneshot::channel();
let (producer_release_tx, producer_release_rx) = oneshot::channel();
let (writer_cancelled_tx, mut writer_cancelled_rx) = oneshot::channel();
let producer = CriticalTask::new(CriticalTaskName::Worker, async move {
producer_shutdown.cancelled().await;
let _ = producer_cancelled_tx.send(());
let _ = producer_release_rx.await;
Ok::<(), CriticalTaskFailure>(())
});
let writer = CriticalTask::new(CriticalTaskName::DatabaseWriter, async move {
writer_shutdown.cancelled().await;
let _ = writer_cancelled_tx.send(());
Ok::<(), CriticalTaskFailure>(())
});
let application = ApplicationRuntime::with_database_writer(
readiness.clone(),
cancellation,
database_shutdown,
Box::pin(async move {
let _ = shutdown_rx.await;
Ok(())
}),
vec![producer],
spawn_critical_task(writer),
);
let running = tokio::spawn(application.run_until_stop());
while !readiness.is_ready() {
tokio::task::yield_now().await;
}
assert!(shutdown_tx.send(()).is_ok());
assert!(producer_cancelled_rx.await.is_ok());
assert!(matches!(
writer_cancelled_rx.try_recv(),
Err(oneshot::error::TryRecvError::Empty)
));
assert!(producer_release_tx.send(()).is_ok());
assert!(writer_cancelled_rx.await.is_ok());
assert!(running.await.unwrap().is_ok());
}
#[tokio::test]
async fn unexpected_success_cancels_and_drains_before_returning_failure() {
let (result, readiness, cancellation, companion_drained) =
run_with_unexpected_task(async { Ok(()) }).await;
assert!(matches!(
result,
Err(ApplicationError::CriticalTaskExited {
task: CriticalTaskName::Scheduler
})
));
assert_shutdown_state(readiness, cancellation, companion_drained);
}
#[tokio::test]
async fn task_error_cancels_and_drains_before_returning_failure() {
let (result, readiness, cancellation, companion_drained) =
run_with_unexpected_task(async {
Err(Box::new(std::io::Error::other("task failed")) as CriticalTaskFailure)
})
.await;
assert!(matches!(
result,
Err(ApplicationError::CriticalTaskFailed {
task: CriticalTaskName::Scheduler,
..
})
));
assert_shutdown_state(readiness, cancellation, companion_drained);
}
#[tokio::test]
async fn task_panic_cancels_and_drains_before_returning_failure() {
async fn panicking_task() -> Result<(), CriticalTaskFailure> {
panic!("task panicked")
}
let (result, readiness, cancellation, companion_drained) =
run_with_unexpected_task(panicking_task()).await;
assert!(matches!(
result,
Err(ApplicationError::CriticalTaskPanicked {
task: CriticalTaskName::Scheduler,
..
})
));
assert_shutdown_state(readiness, cancellation, companion_drained);
}
#[tokio::test]
async fn metrics_collector_storage_failure_marks_unready_and_drains_the_writer() {
let (root, config_path, _publication_path) = startup_fixture("", VALID_PUBLICATION);
let startup =
StartupConfiguration::load_with_discovery(config_path, discover_content_tree).unwrap();
let host = startup._host.view();
let bootstrapped = database::bootstrap(host.database).await.unwrap();
let (store, writer) = bootstrapped.into_store(4);
let metrics = Metrics::new(&store.health.metrics).unwrap();
let collector = MetricsCollector::new(
metrics,
store.health.clone(),
BackupHealth::new(None),
tokio::runtime::Handle::current(),
);
fs::rename(root.path().join("state"), root.path().join("moved-state")).unwrap();
fs::write(root.path().join("state"), []).unwrap();
let readiness = Readiness::default();
let cancellation = CancellationToken::new();
let writer_shutdown = CancellationToken::new();
let application = ApplicationRuntime::with_database_writer(
readiness.clone(),
cancellation.clone(),
writer_shutdown.clone(),
Box::pin(std::future::pending()),
vec![CriticalTask::new(
CriticalTaskName::MetricsCollector,
collector.run(cancellation.clone()),
)],
spawn_critical_task(CriticalTask::new(
CriticalTaskName::DatabaseWriter,
writer.run(writer_shutdown.clone()),
)),
);
let result = tokio::time::timeout(
std::time::Duration::from_secs(5),
application.run_until_stop(),
)
.await
.unwrap();
assert!(matches!(
result,
Err(ApplicationError::CriticalTaskFailed {
task: CriticalTaskName::MetricsCollector,
..
})
));
assert!(!readiness.is_ready());
assert!(cancellation.is_cancelled());
assert!(writer_shutdown.is_cancelled());
assert_eq!(store.health.metrics.writer_up.get(), 0);
}
#[tokio::test]
async fn database_writer_failure_marks_unready_and_drains_producers() {
let readiness = Readiness::default();
let cancellation = CancellationToken::new();
let database_shutdown = CancellationToken::new();
let companion_drained = Arc::new(AtomicBool::new(false));
let companion = {
let cancellation = cancellation.clone();
let companion_drained = Arc::clone(&companion_drained);
CriticalTask::new(CriticalTaskName::Worker, async move {
cancellation.cancelled().await;
companion_drained.store(true, Ordering::SeqCst);
Ok::<(), CriticalTaskFailure>(())
})
};
let writer = CriticalTask::new(CriticalTaskName::DatabaseWriter, async {
Err(std::io::Error::other("writer stopped"))
});
let application = ApplicationRuntime::with_database_writer(
readiness.clone(),
cancellation.clone(),
database_shutdown.clone(),
Box::pin(std::future::pending()),
vec![companion],
spawn_critical_task(writer),
);
let result = application.run_until_stop().await;
assert!(matches!(
result,
Err(ApplicationError::CriticalTaskFailed {
task: CriticalTaskName::DatabaseWriter,
source,
}) if source.downcast_ref::<std::io::Error>()
.is_some_and(|error| error.to_string() == "writer stopped")
));
assert!(!readiness.is_ready());
assert!(cancellation.is_cancelled());
assert!(database_shutdown.is_cancelled());
assert!(companion_drained.load(Ordering::SeqCst));
}
async fn run_with_unexpected_task<Future>(
trigger: Future,
) -> (
Result<(), ApplicationError>,
Readiness,
CancellationToken,
Arc<AtomicBool>,
)
where
Future: std::future::Future<Output = CriticalTaskResult> + Send + 'static,
{
let readiness = Readiness::default();
let cancellation = CancellationToken::new();
let companion_drained = Arc::new(AtomicBool::new(false));
let companion = {
let cancellation = cancellation.clone();
let companion_drained = Arc::clone(&companion_drained);
async move {
cancellation.cancelled().await;
companion_drained.store(true, Ordering::SeqCst);
Ok::<(), CriticalTaskFailure>(())
}
};
let application = ApplicationRuntime::with_parts(
readiness.clone(),
cancellation.clone(),
Box::pin(std::future::pending()),
vec![
CriticalTask::new(CriticalTaskName::Scheduler, trigger),
CriticalTask::new(CriticalTaskName::Worker, companion),
],
);
let result = application.run_until_stop().await;
(result, readiness, cancellation, companion_drained)
}
fn assert_shutdown_state(
readiness: Readiness,
cancellation: CancellationToken,
companion_drained: Arc<AtomicBool>,
) {
assert!(!readiness.is_ready());
assert!(cancellation.is_cancelled());
assert!(companion_drained.load(Ordering::SeqCst));
}
}