1use std::{
2 collections::BTreeMap,
3 error::Error as StdError,
4 future::Future,
5 panic::{AssertUnwindSafe, catch_unwind},
6 path::{Path, PathBuf},
7 pin::Pin,
8 sync::Arc,
9 task::{Context, Poll},
10};
11
12use clap::error::ErrorKind;
13use markdown_compiler::{
14 ContentCandidateStore, ContentCandidateStoreError, ContentTreeDigest, ContentTreeLimits,
15 ContentValidationErrors, DiscoveredContentTree, PostId, PostRevisionDigest,
16 PrepareContentError, PreparedContent, SiteSnapshotDigest, discover_content_tree,
17 prepare_content,
18};
19use time::OffsetDateTime;
20use tokio::task::{JoinError, JoinHandle, JoinSet};
21use tokio_util::sync::CancellationToken;
22use tracing::Instrument as _;
23
24use crate::{
25 admin::{
26 AdminSecurityState, AdminServer, AdminSessionPolicy, origin::AdminBind,
27 runtime_admin_router,
28 },
29 backup_health::BackupHealth,
30 cli::{ServerInvocation, parse_process_invocation},
31 config::{HostConfiguration, HostConfigurationLoader, SourceConfigurationView},
32 content_sync::ContentSync,
33 database::{self, DatabaseStore},
34 domain::{
35 auth::{Argon2idPolicy, store::ConfiguredLoginProviders},
36 mail::{
37 config::MailConfiguration,
38 dispatch::MailDispatcher,
39 feedback::FeedbackWorker,
40 retention::MailRetention,
41 runtime::{PreparedMail, prepare_mail},
42 ui::MailUiState,
43 },
44 profile::TipRecipientProjection,
45 publication::{
46 PublicLedgerProjection, SourceCommit,
47 activation::{
48 PreparedPublicationRecovery, PublicationActivationError, PublicationCoordinator,
49 PublicationCoordinatorActor, PublicationCoordinatorHandle, observed_post_revisions,
50 },
51 scheduler::PublicationScheduler,
52 store::{
53 InstallStartupSnapshot, ObservedPostRevision, RecoverablePublicationActivation,
54 RetainedReleaseInput, StartupSnapshotState,
55 },
56 },
57 },
58 error::{
59 ApplicationError, CriticalTaskName, ProcessError, ProcessExit, ShutdownSignal, StartupStage,
60 },
61 frontend_assets::{FrontendAssetManifest, embedded_manifest},
62 git_sync::GitSync,
63 identity_bootstrap,
64 metrics::{Metrics, MetricsCollector, MetricsServer},
65 observability::{initialize_logging, task_span},
66 process_lock::{ProcessLock, ProcessLockError},
67 render::{
68 CatalogBuildError, CatalogRetentionError, ContentCatalog, ContentCompiler, SiteSnapshot,
69 render_site_shell, snapshot_store,
70 },
71 restore,
72 source_bootstrap::{configure_source, generate_source_key},
73 source_provenance::{SourceCommitDiscovery, discover_source_commit},
74 source_sync::{ManagedSourceEngine, SourceSyncHandle},
75 web::{PublicServer, PublicState, Readiness, public_router_with_routes},
76};
77
78#[cfg(test)]
79use crate::{domain::publication::PublishedPostRevision, render::compile_content_catalog};
80
81type ShutdownFuture = Pin<Box<dyn Future<Output = Result<(), ApplicationError>> + Send>>;
82type CriticalTaskFuture = Pin<Box<dyn Future<Output = CriticalTaskResult> + Send>>;
83type CriticalTaskResult = Result<(), CriticalTaskFailure>;
84type CriticalTaskFailure = Box<dyn std::error::Error + Send + Sync>;
85const PUBLICATION_COORDINATOR_QUEUE_CAPACITY: usize = 32;
86
87pub(crate) struct Application {
93 _process_lock: ProcessLock,
94 _database: DatabaseStore,
95 publication_coordinator: PublicationCoordinatorHandle,
96 runtime: ApplicationRuntime,
97 #[cfg(test)]
98 public_addr: std::net::SocketAddr,
99 #[cfg(test)]
100 admin_addr: std::net::SocketAddr,
101}
102
103struct ApplicationRuntime {
104 readiness: Readiness,
105 cancellation: CancellationToken,
106 database_shutdown: CancellationToken,
107 shutdown: ShutdownFuture,
108 critical_tasks: JoinSet<(CriticalTaskName, CriticalTaskCompletion)>,
109 database_writer: Option<JoinHandle<(CriticalTaskName, CriticalTaskCompletion)>>,
110}
111
112struct StartupConfiguration {
113 _process_lock: ProcessLock,
114 _host: HostConfiguration,
115 _content_tree: DiscoveredContentTree,
116 content: PreparedContent,
117}
118
119struct StartupHostConfiguration {
125 process_lock: ProcessLock,
126 host: HostConfiguration,
127}
128
129struct CompiledStartupContent {
130 catalog: Arc<ContentCatalog>,
131 observed_posts: Vec<ObservedPostRevision>,
132 source_commit: Option<SourceCommit>,
133 content_digest: ContentTreeDigest,
134}
135
136struct ServingState {
137 readiness: Readiness,
138 publication_coordinator: PublicationCoordinatorHandle,
139 publication_actor: PublicationCoordinatorActor,
140 public_server: PublicServer,
141 admin_server: AdminServer,
142 metrics_server: MetricsServer,
143 metrics_collector: MetricsCollector,
144 mail_dispatcher: Option<MailDispatcher>,
145 mail_feedback: Option<FeedbackWorker>,
146}
147
148struct StartedDatabase {
149 store: DatabaseStore,
150 shutdown: CancellationToken,
151 task: JoinHandle<(CriticalTaskName, CriticalTaskCompletion)>,
152 security: AdminSecurityState,
153}
154
155impl StartupConfiguration {
156 #[cfg(test)]
157 fn load_with_discovery<Discover>(
158 config_path: PathBuf,
159 discover: Discover,
160 ) -> Result<Self, ProcessError>
161 where
162 Discover: FnOnce(
163 &Path,
164 ContentTreeLimits,
165 ) -> Result<DiscoveredContentTree, ContentValidationErrors>,
166 {
167 StartupHostConfiguration::load(config_path)?.discover_with(discover)
168 }
169}
170
171impl StartupHostConfiguration {
172 fn load(config_path: PathBuf) -> Result<Self, ProcessError> {
173 let host = HostConfigurationLoader::from_process_working_directory()?.load(&config_path)?;
174 let host_view = host.view();
175 let process_lock = match ProcessLock::acquire(host_view.runtime_root) {
176 Ok(process_lock) => process_lock,
177 Err(ProcessLockError::AlreadyRunning) => return Err(ProcessError::AlreadyRunning),
178 Err(error) => {
179 return Err(startup_failure(
180 StartupStage::ProcessLock,
181 "acquire the process lock",
182 error,
183 ));
184 }
185 };
186
187 Ok(Self { process_lock, host })
188 }
189
190 fn discover_with<Discover>(
191 self,
192 discover: Discover,
193 ) -> Result<StartupConfiguration, ProcessError>
194 where
195 Discover: FnOnce(
196 &Path,
197 ContentTreeLimits,
198 ) -> Result<DiscoveredContentTree, ContentValidationErrors>,
199 {
200 let host_view = self.host.view();
201 let content_tree = discover(host_view.content_root, host_view.content_limits)?;
202 let content = prepare_content(&content_tree).map_err(|error| match error {
203 PrepareContentError::InvalidContent(source) => ProcessError::Validation(source),
204 PrepareContentError::AssetResolution(source) => {
205 startup_failure(StartupStage::Content, "resolve content assets", source)
206 }
207 })?;
208
209 Ok(StartupConfiguration {
210 _process_lock: self.process_lock,
211 _host: self.host,
212 _content_tree: content_tree,
213 content,
214 })
215 }
216}
217
218pub async fn run_until_stop() -> ProcessExit {
220 initialize_logging();
221 let invocation = match parse_process_invocation() {
222 Ok(invocation) => invocation,
223 Err(error) => return report_command_error(error),
224 };
225 let result: Result<(), ProcessError> = async {
226 match invocation {
227 ServerInvocation::CheckpointManifest {
228 config_path,
229 database_file,
230 plan_file,
231 ltx_root,
232 artifact_root,
233 output,
234 } => {
235 restore::checkpoint::write_manifest(
236 config_path,
237 database_file,
238 plan_file,
239 ltx_root,
240 artifact_root,
241 output,
242 )
243 .await
244 }
245 ServerInvocation::VerifyCheckpoint {
246 manifest_file,
247 ltx_root,
248 artifact_root,
249 } => {
250 restore::checkpoint::verify_checkpoint(manifest_file, ltx_root, artifact_root).await
251 }
252 ServerInvocation::RestoreReplica {
253 config_path,
254 database_file,
255 artifact_root,
256 manifest_file,
257 ltx_root,
258 } => {
259 restore::restore(
260 config_path,
261 database_file,
262 artifact_root,
263 restore::RestoreManifest::Replica {
264 path: manifest_file,
265 ltx_root,
266 },
267 )
268 .await
269 }
270 ServerInvocation::ExportBackup {
271 config_path,
272 database_file,
273 } => restore::export_backup(config_path, database_file).await,
274 ServerInvocation::Restore {
275 config_path,
276 database_file,
277 artifact_root,
278 manifest_file,
279 } => {
280 restore::restore(
281 config_path,
282 database_file,
283 artifact_root,
284 restore::RestoreManifest::Bundle(manifest_file),
285 )
286 .await
287 }
288 ServerInvocation::Serve { config_path } => {
289 let startup = StartupHostConfiguration::load(config_path)?;
290 let application = match startup.host.view().source {
291 SourceConfigurationView::ExternalCheckout => {
292 Application::build(startup.discover_with(discover_content_tree)?).await?
293 }
294 SourceConfigurationView::ManagedGit { .. } => {
295 Application::build_managed(startup).await?
296 }
297 };
298 application
299 .run_until_stop()
300 .await
301 .map_err(ProcessError::from)
302 }
303 ServerInvocation::BootstrapIdentity {
304 config_path,
305 credential,
306 } => identity_bootstrap::bootstrap_owner(config_path, credential).await,
307 ServerInvocation::ConfigureSource {
308 config_path,
309 request,
310 idempotency_key,
311 } => configure_source(config_path, request, idempotency_key).await,
312 ServerInvocation::GenerateSourceKey {
313 config_path,
314 private_key_file,
315 } => generate_source_key(config_path, private_key_file).await,
316 }
317 }
318 .await;
319
320 match result {
321 Ok(()) => ProcessExit::Success,
322 Err(error) => {
323 let exit = error.exit();
324 tracing::error!(
325 error = %error,
326 category = error.category(),
327 exit_code = exit.code(),
328 "server process failed"
329 );
330 exit
331 }
332 }
333}
334
335fn report_command_error(error: clap::Error) -> ProcessExit {
336 let exit = match error.kind() {
337 ErrorKind::DisplayHelp | ErrorKind::DisplayVersion => ProcessExit::Success,
338 _ => ProcessExit::Usage,
339 };
340
341 if let Err(print_error) = error.print() {
342 tracing::error!(
343 error = %print_error,
344 "failed to print command output"
345 );
346 return ProcessExit::Internal;
347 }
348
349 exit
350}
351
352impl Application {
353 async fn build(startup: StartupConfiguration) -> Result<Self, ProcessError> {
354 let cancellation = CancellationToken::new();
355 let host = startup._host.view();
356 let content_root = host.content_root.to_path_buf();
357 let state_root = host.state_root.to_path_buf();
358 let content_limits = host.content_limits;
359 let frontend = embedded_manifest();
360 frontend.validate().map_err(|error| {
361 startup_failure(
362 StartupStage::FrontendAssets,
363 "validate embedded frontend assets",
364 error,
365 )
366 })?;
367 let shutdown = install_termination_signal()?;
368 let candidate_store =
369 ContentCandidateStore::open(&state_root, content_limits).map_err(|error| {
370 startup_failure(
371 StartupStage::Content,
372 "open the retained content candidate store",
373 error,
374 )
375 })?;
376 candidate_store
377 .retain(&startup._content_tree)
378 .map_err(|error| {
379 startup_failure(
380 StartupStage::Content,
381 "retain the startup content candidate",
382 error,
383 )
384 })?;
385 let content_compiler = ContentCompiler::discover().map_err(|error| {
386 startup_failure(
387 StartupStage::Content,
388 "initialize the content rendering pipeline",
389 error,
390 )
391 })?;
392 let compiled = compile_startup_content(&startup, host.content_root, &content_compiler)?;
393 let active_content_digest = compiled.content_digest.clone();
394 let database = start_database(&startup._host).await?;
395 let source = SourceSyncHandle::external_checkout(database.store.source.clone());
396
397 let serving_state = match prepare_serving_state(ServingStateInput {
398 database: &database.store,
399 compiled,
400 frontend,
401 public_bind: host.public_bind,
402 admin_bind: host.admin_bind,
403 metrics_bind: host.metrics_bind,
404 backup: BackupHealth::new(host.backup),
405 security: database.security.clone(),
406 mail: host.mail,
407 cancellation: cancellation.clone(),
408 candidate_store: &candidate_store,
409 content_compiler: &content_compiler,
410 source,
411 })
412 .await
413 {
414 Ok(setup) => setup,
415 Err(error) => {
416 return Err(close_started_database(database, error).await);
417 }
418 };
419 let content_sync = ContentSync::new(
420 content_root,
421 content_limits,
422 candidate_store,
423 active_content_digest,
424 serving_state.publication_coordinator.clone(),
425 cancellation.clone(),
426 content_compiler,
427 );
428 let content_task = CriticalTask::new(CriticalTaskName::ContentSync, content_sync.run());
429 Ok(Self::assemble(
430 startup._process_lock,
431 database,
432 serving_state,
433 cancellation,
434 shutdown,
435 content_task,
436 ))
437 }
438
439 async fn build_managed(startup: StartupHostConfiguration) -> Result<Self, ProcessError> {
440 let cancellation = CancellationToken::new();
441 let host = startup.host.view();
442 let state_root = host.state_root.to_path_buf();
443 let content_limits = host.content_limits;
444 let (mirror_root, credentials, process_limits) = match host.source {
445 SourceConfigurationView::ManagedGit {
446 mirror_root,
447 credentials,
448 limits,
449 } => (mirror_root, credentials, limits),
450 SourceConfigurationView::ExternalCheckout => {
451 return Err(ProcessError::ManagedSourceDisabled);
452 }
453 };
454 let frontend = embedded_manifest();
455 frontend.validate().map_err(|error| {
456 startup_failure(
457 StartupStage::FrontendAssets,
458 "validate embedded frontend assets",
459 error,
460 )
461 })?;
462 let shutdown = install_termination_signal()?;
463 let database = start_database(&startup.host).await?;
464 match database.store.source.configuration().await {
465 Ok(Some(_)) => {}
466 Ok(None) => {
467 return Err(close_started_database(
468 database,
469 ProcessError::SourceConfigurationRequired,
470 )
471 .await);
472 }
473 Err(error) => {
474 let error = startup_failure(
475 StartupStage::Source,
476 "load durable managed-source settings",
477 error,
478 );
479 return Err(close_started_database(database, error).await);
480 }
481 };
482 let git = match GitSync::discover(mirror_root, credentials, process_limits, content_limits)
483 {
484 Ok(git) => git,
485 Err(error) => {
486 let error = startup_failure(
487 StartupStage::Source,
488 "initialize the managed Git transport",
489 error,
490 );
491 return Err(close_started_database(database, error).await);
492 }
493 };
494 let candidate_store = match ContentCandidateStore::open(&state_root, content_limits) {
495 Ok(store) => store,
496 Err(error) => {
497 let error = startup_failure(
498 StartupStage::Content,
499 "open the retained content candidate store",
500 error,
501 );
502 return Err(close_started_database(database, error).await);
503 }
504 };
505 let content_compiler = match ContentCompiler::discover() {
506 Ok(compiler) => compiler,
507 Err(error) => {
508 let error = startup_failure(
509 StartupStage::Content,
510 "initialize the content rendering pipeline",
511 error,
512 );
513 return Err(close_started_database(database, error).await);
514 }
515 };
516 let (source_engine, source) = ManagedSourceEngine::new(
517 database.store.source.clone(),
518 git,
519 candidate_store.clone(),
520 content_compiler.clone(),
521 cancellation.clone(),
522 );
523 let candidate = match source_engine.prepare_startup().await {
524 Ok(candidate) => candidate,
525 Err(error) => {
526 let error = startup_failure(
527 StartupStage::Source,
528 "prepare the managed source head",
529 error,
530 );
531 return Err(close_started_database(database, error).await);
532 }
533 };
534 let compiled = CompiledStartupContent {
535 observed_posts: observed_post_revisions(&candidate.catalog),
536 catalog: candidate.catalog,
537 source_commit: Some(candidate.source_commit),
538 content_digest: candidate.content_digest,
539 };
540 let serving_state = match prepare_serving_state(ServingStateInput {
541 database: &database.store,
542 compiled,
543 frontend,
544 public_bind: host.public_bind,
545 admin_bind: host.admin_bind,
546 metrics_bind: host.metrics_bind,
547 backup: BackupHealth::new(host.backup),
548 security: database.security.clone(),
549 mail: host.mail,
550 cancellation: cancellation.clone(),
551 candidate_store: &candidate_store,
552 content_compiler: &content_compiler,
553 source,
554 })
555 .await
556 {
557 Ok(setup) => setup,
558 Err(error) => {
559 return Err(close_started_database(database, error).await);
560 }
561 };
562 let source_sync = source_engine.into_live(serving_state.publication_coordinator.clone());
563 let source_task = CriticalTask::new(CriticalTaskName::SourceSync, source_sync.run());
564 Ok(Self::assemble(
565 startup.process_lock,
566 database,
567 serving_state,
568 cancellation,
569 shutdown,
570 source_task,
571 ))
572 }
573
574 fn assemble(
576 process_lock: ProcessLock,
577 database: StartedDatabase,
578 serving_state: ServingState,
579 cancellation: CancellationToken,
580 shutdown: ShutdownFuture,
581 source_task: CriticalTask,
582 ) -> Self {
583 let ServingState {
584 readiness,
585 publication_coordinator,
586 publication_actor,
587 public_server,
588 admin_server,
589 metrics_server,
590 metrics_collector,
591 mail_dispatcher,
592 mail_feedback,
593 } = serving_state;
594 #[cfg(test)]
595 let public_addr = public_server.local_addr;
596 #[cfg(test)]
597 let admin_addr = admin_server.local_addr;
598 let public_task = CriticalTask::new(
599 CriticalTaskName::PublicServer,
600 public_server.serve(cancellation.clone()),
601 );
602 let admin_task = CriticalTask::new(
603 CriticalTaskName::AdminServer,
604 admin_server.serve(cancellation.clone()),
605 );
606 let metrics_task = CriticalTask::new(
607 CriticalTaskName::MetricsServer,
608 metrics_server.serve(cancellation.clone()),
609 );
610 let collector_task = CriticalTask::new(
611 CriticalTaskName::MetricsCollector,
612 metrics_collector.run(cancellation.clone()),
613 );
614 let publication_actor_task = CriticalTask::new(
615 CriticalTaskName::PublicationCoordinator,
616 publication_actor.run(cancellation.clone()),
617 );
618 let scheduler = PublicationScheduler::new(
619 database.store.publications.clone(),
620 publication_coordinator.clone(),
621 publication_coordinator.scheduler_wakeup(),
622 cancellation.clone(),
623 );
624 let scheduler_task = CriticalTask::new(CriticalTaskName::Scheduler, scheduler.run());
625
626 let retention = MailRetention::new(database.store.subscribers.clone());
627 let retention_task = CriticalTask::new(
628 CriticalTaskName::MailRetention,
629 retention.run(cancellation.clone()),
630 );
631 let mut tasks = vec![
632 retention_task,
633 publication_actor_task,
634 public_task,
635 admin_task,
636 metrics_task,
637 collector_task,
638 source_task,
639 scheduler_task,
640 ];
641 if let Some(dispatcher) = mail_dispatcher {
642 tasks.push(CriticalTask::new(
643 CriticalTaskName::MailDispatch,
644 dispatcher.run(cancellation.clone()),
645 ));
646 }
647 if let Some(feedback) = mail_feedback {
648 tasks.push(CriticalTask::new(
649 CriticalTaskName::MailFeedback,
650 feedback.run(cancellation.clone()),
651 ));
652 }
653 Self {
654 _process_lock: process_lock,
655 _database: database.store,
656 publication_coordinator,
657 runtime: ApplicationRuntime::with_database_writer(
658 readiness,
659 cancellation,
660 database.shutdown,
661 shutdown,
662 tasks,
663 database.task,
664 ),
665 #[cfg(test)]
666 public_addr,
667 #[cfg(test)]
668 admin_addr,
669 }
670 }
671
672 async fn run_until_stop(self) -> Result<(), ApplicationError> {
673 let Self {
674 _process_lock: process_lock,
675 _database: database,
676 publication_coordinator,
677 runtime,
678 #[cfg(test)]
679 public_addr: _,
680 #[cfg(test)]
681 admin_addr: _,
682 } = self;
683 let runtime_result = runtime.run_until_stop().await;
684 drop(publication_coordinator);
685 drop(database);
686 drop(process_lock);
687 runtime_result
688 }
689}
690
691fn compile_startup_content(
692 startup: &StartupConfiguration,
693 content_root: &Path,
694 compiler: &ContentCompiler,
695) -> Result<CompiledStartupContent, ProcessError> {
696 let content_digest = startup._content_tree.digest();
697 let catalog = Arc::new(compiler.compile(&startup.content).map_err(|error| {
698 startup_failure(StartupStage::Content, "compile the content catalog", error)
699 })?);
700 let observed_posts = observed_post_revisions(&catalog);
701 let source_commit = match discover_source_commit(content_root) {
702 SourceCommitDiscovery::Discovered(commit) => Some(commit),
703 SourceCommitDiscovery::Unavailable(reason) => {
704 tracing::warn!(?reason, "content source commit is unavailable");
705 None
706 }
707 };
708 Ok(CompiledStartupContent {
709 catalog,
710 observed_posts,
711 source_commit,
712 content_digest,
713 })
714}
715
716pub(crate) async fn verify_restore_content(
718 database: &DatabaseStore,
719 state_root: &Path,
720 limits: ContentTreeLimits,
721) -> Result<(), ProcessError> {
722 let inputs = RestoreContentInputs::load(database).await?;
723 let state_root = state_root.to_path_buf();
724 tokio::task::spawn_blocking(move || inputs.verify(&state_root, limits))
725 .await
726 .map_err(|error| {
727 startup_failure(
728 StartupStage::Content,
729 "await offline restore verification",
730 error,
731 )
732 })?
733}
734
735struct RestoreContentInputs {
737 tip_recipient: Option<TipRecipientProjection>,
738 startup: StartupSnapshotState,
739 installed_content: Option<ContentTreeDigest>,
740 retained: Vec<RetainedReleaseInput>,
741}
742
743impl RestoreContentInputs {
744 async fn load(database: &DatabaseStore) -> Result<Self, ProcessError> {
745 let tip_recipient = database
746 .profiles
747 .effective_tip_recipient()
748 .await
749 .map_err(|error| {
750 startup_failure(
751 StartupStage::Database,
752 "verify restored tip projection",
753 error,
754 )
755 })?;
756 let startup = database
757 .publications
758 .startup_snapshot_state()
759 .await
760 .map_err(|error| {
761 startup_failure(
762 StartupStage::Database,
763 "verify restored publication ledger",
764 error,
765 )
766 })?;
767 let source = database.source.status().await.map_err(|error| {
768 startup_failure(
769 StartupStage::Database,
770 "verify restored source ledger",
771 error,
772 )
773 })?;
774 let retained = database
775 .publications
776 .retained_release_inputs()
777 .await
778 .map_err(|error| {
779 startup_failure(
780 StartupStage::Database,
781 "verify retained restore approvals",
782 error,
783 )
784 })?;
785 Ok(Self {
786 tip_recipient,
787 startup,
788 installed_content: source.installation.map(|source| source.content_digest),
789 retained,
790 })
791 }
792
793 fn verify(&self, state_root: &Path, limits: ContentTreeLimits) -> Result<(), ProcessError> {
794 let candidates = ContentCandidateStore::open(state_root, limits).map_err(|error| {
795 startup_failure(
796 StartupStage::Content,
797 "inspect restored candidate archives",
798 error,
799 )
800 })?;
801 let compiler = ContentCompiler::discover().map_err(|error| {
802 startup_failure(StartupStage::Content, "initialize restore compiler", error)
803 })?;
804 let catalogs = compile_retained_catalogs(&candidates, &compiler).map_err(|error| {
805 startup_failure(
806 StartupStage::Content,
807 "compile restored candidate archives",
808 error,
809 )
810 })?;
811 let base = catalogs.values().next().ok_or_else(|| {
812 startup_failure(
813 StartupStage::Content,
814 "verify retained restore inputs",
815 RetainedCatalogError::CandidateUnavailable,
816 )
817 })?;
818 let pins = self.startup.ledger.revision_keys().chain(
819 self.retained
820 .iter()
821 .map(|input| (input.post_id.clone(), input.revision.clone())),
822 );
823 let preview = hydrate_catalog(base.as_ref().clone(), &catalogs, pins).map_err(|error| {
824 startup_failure(
825 StartupStage::Content,
826 "verify restored release revisions",
827 error,
828 )
829 })?;
830 for digest in self
831 .installed_content
832 .iter()
833 .chain(self.retained.iter().map(|item| &item.content_digest))
834 {
835 if !catalogs.contains_key(digest) {
836 return Err(startup_failure(
837 StartupStage::Content,
838 "verify restored pinned candidates",
839 RetainedCatalogError::CandidateUnavailable,
840 ));
841 }
842 }
843 self.verify_activations(&catalogs, &preview)?;
844 rebuild_public_snapshot(
845 &self.startup,
846 &catalogs,
847 &preview,
848 embedded_manifest(),
849 self.tip_recipient.as_ref(),
850 )?;
851 Ok(())
852 }
853
854 fn verify_activations(
856 &self,
857 catalogs: &BTreeMap<ContentTreeDigest, Arc<ContentCatalog>>,
858 preview: &Arc<ContentCatalog>,
859 ) -> Result<(), ProcessError> {
860 for activation in self.startup.activating.iter().cloned() {
861 rebuild_activating_publication(
862 activation,
863 &self.startup,
864 catalogs,
865 preview,
866 embedded_manifest(),
867 self.tip_recipient.as_ref(),
868 )
869 .map_err(|error| {
870 startup_failure(StartupStage::Content, "verify restored activation", error)
871 })?;
872 }
873 Ok(())
874 }
875}
876
877#[derive(Debug, thiserror::Error)]
878enum StartupRecoveryError {
879 #[error("the activating publication has no durable site head")]
880 MissingSite,
881 #[error("the activating publication's retained content candidate is unavailable")]
882 MissingCandidate,
883 #[error("the activating publication's retained revisions could not be hydrated")]
884 Revisions(#[from] CatalogRetentionError),
885 #[error("the activating publication could not be reconstructed")]
886 Activation(#[from] PublicationActivationError),
887}
888
889fn rebuild_activating_publication(
891 activation: RecoverablePublicationActivation,
892 startup: &StartupSnapshotState,
893 catalogs: &BTreeMap<ContentTreeDigest, Arc<ContentCatalog>>,
894 preview: &ContentCatalog,
895 frontend: &'static FrontendAssetManifest,
896 tip_recipient: Option<&TipRecipientProjection>,
897) -> Result<PreparedPublicationRecovery, StartupRecoveryError> {
898 let site = startup
899 .site
900 .clone()
901 .ok_or(StartupRecoveryError::MissingSite)?;
902 let mut catalog = catalogs
903 .get(&activation.content_digest)
904 .ok_or(StartupRecoveryError::MissingCandidate)?
905 .as_ref()
906 .clone();
907 catalog.retain_revisions_from(preview, startup.ledger.revision_keys())?;
908 Ok(PreparedPublicationRecovery::prepare(
909 activation,
910 Arc::new(catalog),
911 frontend,
912 &startup.ledger,
913 tip_recipient,
914 site,
915 )?)
916}
917
918fn compile_retained_catalogs(
919 store: &ContentCandidateStore,
920 compiler: &ContentCompiler,
921) -> Result<BTreeMap<ContentTreeDigest, Arc<ContentCatalog>>, RetainedCatalogError> {
922 store
923 .load_all()
924 .map_err(RetainedCatalogError::Load)?
925 .into_iter()
926 .map(|candidate| {
927 let digest = candidate.digest;
928 let content = prepare_content(&candidate.tree).map_err(|source| {
929 RetainedCatalogError::Prepare {
930 digest: digest.clone(),
931 source,
932 }
933 })?;
934 match compiler.compile(&content) {
935 Ok(catalog) => Ok((digest, Arc::new(catalog))),
936 Err(source) => Err(RetainedCatalogError::Compile { digest, source }),
937 }
938 })
939 .collect()
940}
941
942fn hydrate_catalog(
943 mut base: ContentCatalog,
944 retained: &BTreeMap<ContentTreeDigest, Arc<ContentCatalog>>,
945 pins: impl IntoIterator<Item = (PostId, PostRevisionDigest)>,
946) -> Result<Arc<ContentCatalog>, RetainedCatalogError> {
947 for (post_id, revision) in pins {
948 if base.get(&post_id, &revision).is_some() {
949 continue;
950 }
951 let source = retained
952 .values()
953 .find(|catalog| catalog.get(&post_id, &revision).is_some())
954 .ok_or_else(|| RetainedCatalogError::RevisionUnavailable {
955 post_id: post_id.clone(),
956 revision: revision.clone(),
957 })?;
958 base.retain_revisions_from(source, std::iter::once((post_id.clone(), revision.clone())))
959 .map_err(RetainedCatalogError::Retain)?;
960 }
961 Ok(Arc::new(base))
962}
963
964fn find_retained_public_snapshot(
965 retained: &BTreeMap<ContentTreeDigest, Arc<ContentCatalog>>,
966 ledger: &PublicLedgerProjection,
967 expected: &SiteSnapshotDigest,
968 frontend: &'static FrontendAssetManifest,
969 tip_recipient: Option<&TipRecipientProjection>,
970) -> Result<SiteSnapshot, RetainedCatalogError> {
971 let pins = ledger
972 .published_posts()
973 .map(|published| (published.post_id.clone(), published.revision.clone()))
974 .collect::<Vec<_>>();
975 for base in retained.values() {
976 let catalog = hydrate_catalog(base.as_ref().clone(), retained, pins.clone())?;
977 let Ok(shell) = render_site_shell(catalog, frontend, ledger) else {
978 continue;
979 };
980 let shell = shell.bind_tip_recipient(tip_recipient.cloned());
981 let Ok(snapshot) = shell.into_snapshot() else {
982 continue;
983 };
984 if &snapshot.digest == expected {
985 return Ok(snapshot);
986 }
987 }
988 Err(RetainedCatalogError::PublicSnapshotUnavailable {
989 expected: expected.clone(),
990 })
991}
992
993fn rebuild_public_snapshot(
994 startup: &StartupSnapshotState,
995 retained: &BTreeMap<ContentTreeDigest, Arc<ContentCatalog>>,
996 preview_catalog: &Arc<ContentCatalog>,
997 frontend: &'static FrontendAssetManifest,
998 tip_recipient: Option<&TipRecipientProjection>,
999) -> Result<SiteSnapshot, ProcessError> {
1000 let Some(expected) = startup.site.as_ref() else {
1001 let shell = render_site_shell(Arc::clone(preview_catalog), frontend, &startup.ledger)
1002 .map_err(|error| {
1003 startup_failure(StartupStage::Content, "render the site shell", error)
1004 })?
1005 .bind_tip_recipient(tip_recipient.cloned());
1006 return shell.into_snapshot().map_err(|error| {
1007 startup_failure(StartupStage::Content, "build the site snapshot", error)
1008 });
1009 };
1010 find_retained_public_snapshot(
1011 retained,
1012 &startup.ledger,
1013 &expected.digest,
1014 frontend,
1015 tip_recipient,
1016 )
1017 .map_err(|error| {
1018 startup_failure(
1019 StartupStage::Content,
1020 "rebuild the approved public site snapshot",
1021 error,
1022 )
1023 })
1024}
1025
1026#[derive(Debug, thiserror::Error)]
1027enum RetainedCatalogError {
1028 #[error("a required retained content candidate is unavailable")]
1029 CandidateUnavailable,
1030 #[error("retained content candidates could not be loaded")]
1031 Load(#[source] ContentCandidateStoreError),
1032 #[error("retained content candidate {digest} could not be prepared")]
1033 Prepare {
1034 digest: ContentTreeDigest,
1035 #[source]
1036 source: PrepareContentError,
1037 },
1038 #[error("retained content candidate {digest} could not be compiled")]
1039 Compile {
1040 digest: ContentTreeDigest,
1041 #[source]
1042 source: CatalogBuildError,
1043 },
1044 #[error("retained revision {revision} for post {post_id} is unavailable")]
1045 RevisionUnavailable {
1046 post_id: PostId,
1047 revision: PostRevisionDigest,
1048 },
1049 #[error("a retained revision could not be installed in the recovery catalog")]
1050 Retain(#[source] CatalogRetentionError),
1051 #[error("no retained content candidate rebuilds durable site {expected}")]
1052 PublicSnapshotUnavailable { expected: SiteSnapshotDigest },
1053}
1054
1055struct ServingStateInput<'resources> {
1056 database: &'resources DatabaseStore,
1057 mail: &'resources MailConfiguration,
1058 compiled: CompiledStartupContent,
1059 frontend: &'static FrontendAssetManifest,
1060 public_bind: std::net::SocketAddr,
1061 admin_bind: AdminBind,
1062 metrics_bind: std::net::SocketAddr,
1063 backup: BackupHealth,
1064 security: AdminSecurityState,
1065 cancellation: CancellationToken,
1066 candidate_store: &'resources ContentCandidateStore,
1067 content_compiler: &'resources ContentCompiler,
1068 source: SourceSyncHandle,
1069}
1070
1071async fn prepare_serving_state(input: ServingStateInput<'_>) -> Result<ServingState, ProcessError> {
1072 let ServingStateInput {
1073 database,
1074 mail,
1075 compiled,
1076 frontend,
1077 public_bind,
1078 admin_bind,
1079 metrics_bind,
1080 backup,
1081 security,
1082 cancellation,
1083 candidate_store,
1084 content_compiler,
1085 source,
1086 } = input;
1087 database
1088 .mail
1089 .quarantine_interrupted(OffsetDateTime::now_utc())
1090 .await
1091 .map_err(|error| {
1092 startup_failure(
1093 StartupStage::Database,
1094 "quarantine interrupted mail campaigns",
1095 error,
1096 )
1097 })?;
1098 let PreparedMail {
1099 access: mail_access,
1100 public_routes: mail_routes,
1101 dispatcher: mail_dispatcher,
1102 feedback: mail_feedback,
1103 } = prepare_mail(
1104 mail,
1105 database,
1106 compiled.catalog.publication.site.base_url.clone(),
1107 )
1108 .await
1109 .map_err(|error| startup_failure(StartupStage::Configuration, "prepare mail runtime", error))?;
1110 let tip_recipient = database
1111 .profiles
1112 .effective_tip_recipient()
1113 .await
1114 .map_err(|error| {
1115 startup_failure(
1116 StartupStage::Database,
1117 "load the active tip recipient profile",
1118 error,
1119 )
1120 })?;
1121 let mut startup_state = database
1122 .publications
1123 .startup_snapshot_state()
1124 .await
1125 .map_err(|error| {
1126 startup_failure(
1127 StartupStage::Database,
1128 "load the startup publication ledger",
1129 error,
1130 )
1131 })?;
1132 let ledger = startup_state.ledger.clone();
1133 let retained_catalogs =
1134 compile_retained_catalogs(candidate_store, content_compiler).map_err(|error| {
1135 startup_failure(
1136 StartupStage::Content,
1137 "compile retained content candidates",
1138 error,
1139 )
1140 })?;
1141 let preview_pins = ledger
1142 .published_posts()
1143 .map(|published| (published.post_id.clone(), published.revision.clone()))
1144 .chain(startup_state.scheduled.iter().map(|scheduled| {
1145 let view = scheduled.publication.view();
1146 (view.stable_post_id.clone(), view.pinned_post_digest.clone())
1147 }))
1148 .chain(startup_state.activating.iter().map(|activation| {
1149 let view = activation.publication.view();
1150 (view.stable_post_id.clone(), view.pinned_post_digest.clone())
1151 }))
1152 .collect::<Vec<_>>();
1153 let preview_catalog = hydrate_catalog(
1154 compiled.catalog.as_ref().clone(),
1155 &retained_catalogs,
1156 preview_pins,
1157 )
1158 .map_err(|error| {
1159 startup_failure(
1160 StartupStage::Content,
1161 "hydrate the private preview catalog",
1162 error,
1163 )
1164 })?;
1165 let recovery = startup_state
1166 .activating
1167 .pop()
1168 .map(|activation| {
1169 rebuild_activating_publication(
1170 activation,
1171 &startup_state,
1172 &retained_catalogs,
1173 &preview_catalog,
1174 frontend,
1175 tip_recipient.as_ref(),
1176 )
1177 })
1178 .transpose()
1179 .map_err(|error| {
1180 startup_failure(
1181 StartupStage::Content,
1182 "rebuild the activating publication candidate",
1183 error,
1184 )
1185 })?;
1186 let snapshot = rebuild_public_snapshot(
1187 &startup_state,
1188 &retained_catalogs,
1189 &preview_catalog,
1190 frontend,
1191 tip_recipient.as_ref(),
1192 )?;
1193 let installed_site = database
1194 .publications
1195 .install_startup_snapshot(InstallStartupSnapshot {
1196 expected: startup_state.site.clone(),
1197 candidate_digest: snapshot.digest.clone(),
1198 activated_at: OffsetDateTime::now_utc(),
1199 source_commit: compiled.source_commit.clone(),
1200 posts: compiled.observed_posts,
1201 })
1202 .await
1203 .map_err(|error| {
1204 startup_failure(
1205 StartupStage::Database,
1206 "install the startup site snapshot",
1207 error,
1208 )
1209 })?;
1210
1211 let readiness = Readiness::default();
1212 let (snapshots, activator) = snapshot_store(snapshot);
1213 let mut publication_coordinator = PublicationCoordinator {
1214 catalog: preview_catalog,
1215 content_digest: compiled.content_digest,
1216 candidates: Arc::new(retained_catalogs),
1217 ledger,
1218 site: installed_site,
1219 activator,
1220 store: database.publications.clone(),
1221 profiles: database.profiles.clone(),
1222 tip_recipient,
1223 frontend,
1224 source_commit: compiled.source_commit,
1225 scheduled: startup_state
1226 .scheduled
1227 .drain(..)
1228 .map(|scheduled| (scheduled.publication_id, scheduled))
1229 .collect(),
1230 scheduler_wakeup: Arc::new(tokio::sync::Notify::new()),
1231 readiness: readiness.clone(),
1232 cancellation: cancellation.clone(),
1233 };
1234 if let Some(recovery) = recovery {
1235 publication_coordinator
1236 .recover(recovery)
1237 .await
1238 .map_err(|error| {
1239 startup_failure(
1240 StartupStage::Database,
1241 "recover the activating publication",
1242 error,
1243 )
1244 })?;
1245 }
1246 let (publication_coordinator, publication_actor) =
1247 publication_coordinator.into_actor(PUBLICATION_COORDINATOR_QUEUE_CAPACITY);
1248 let protected_admin_router = runtime_admin_router(
1249 publication_coordinator.clone(),
1250 security,
1251 database.profiles.clone(),
1252 source,
1253 MailUiState {
1254 campaigns: database.mail.clone(),
1255 subscribers: database.subscribers.clone(),
1256 publications: publication_coordinator.clone(),
1257 snapshots: snapshots.clone(),
1258 access: mail_access,
1259 },
1260 );
1261 let public_state = PublicState {
1262 snapshots,
1263 readiness: readiness.clone(),
1264 };
1265 let public_server = match mail_routes {
1266 Some(routes) => {
1267 PublicServer::bind_router(public_bind, public_router_with_routes(public_state, routes))
1268 .await
1269 }
1270 None => PublicServer::bind(public_bind, public_state).await,
1271 }
1272 .map_err(|error| startup_failure(StartupStage::Listeners, "bind the public listener", error))?;
1273 tracing::info!(bind = %public_server.local_addr, "public listener bound");
1274 let admin_server = AdminServer::bind(admin_bind, protected_admin_router)
1275 .await
1276 .map_err(|error| {
1277 startup_failure(StartupStage::Listeners, "bind the admin listener", error)
1278 })?;
1279 tracing::info!(
1280 bind = %admin_server.local_addr,
1281 "authenticated admin backend listener bound"
1282 );
1283 let metrics = Metrics::new(&database.health.metrics).map_err(|error| {
1284 startup_failure(
1285 StartupStage::Listeners,
1286 "construct the application metrics registry",
1287 error,
1288 )
1289 })?;
1290 let metrics_server = MetricsServer::bind(metrics_bind, metrics.clone())
1291 .await
1292 .map_err(|error| {
1293 startup_failure(StartupStage::Listeners, "bind the metrics listener", error)
1294 })?;
1295 tracing::info!(bind = %metrics_server.local_addr, "loopback metrics listener bound");
1296 let metrics_collector = MetricsCollector::new(
1297 metrics,
1298 database.health.clone(),
1299 backup,
1300 tokio::runtime::Handle::current(),
1301 );
1302 Ok(ServingState {
1303 mail_dispatcher,
1304 mail_feedback,
1305 metrics_server,
1306 metrics_collector,
1307 readiness,
1308 publication_coordinator,
1309 publication_actor,
1310 public_server,
1311 admin_server,
1312 })
1313}
1314
1315fn startup_failure(
1316 stage: StartupStage,
1317 operation: &'static str,
1318 error: impl StdError + Send + Sync + 'static,
1319) -> ProcessError {
1320 tracing::error!(%stage, operation, error = %error, "startup operation failed");
1321 ApplicationError::Startup {
1322 stage,
1323 operation,
1324 source: Box::new(error),
1325 }
1326 .into()
1327}
1328
1329async fn start_database(host: &HostConfiguration) -> Result<StartedDatabase, ProcessError> {
1330 let host = host.view();
1331 restore::verify_startup_candidate(host.database.path, host.state_root)
1332 .await
1333 .map_err(|error| {
1334 startup_failure(
1335 StartupStage::Database,
1336 "verify restored startup candidate",
1337 error,
1338 )
1339 })?;
1340 let database = match database::bootstrap(host.database).await {
1341 Ok(database) => database,
1342 Err(database::DatabaseStartupError::AlreadyOwned) => {
1343 return Err(ProcessError::AlreadyRunning);
1344 }
1345 Err(error) => {
1346 return Err(startup_failure(
1347 StartupStage::Database,
1348 "bootstrap the database",
1349 error,
1350 ));
1351 }
1352 };
1353 let (store, writer) = database.into_store(host.database.writer_queue_capacity.get());
1354 let shutdown = CancellationToken::new();
1355 let task = spawn_critical_task(CriticalTask::new(
1356 CriticalTaskName::DatabaseWriter,
1357 writer.run(shutdown.clone()),
1358 ));
1359
1360 if let Err(error) = identity_bootstrap::initialize_startup_identity(
1361 &store,
1362 host.identity_startup_bootstrap,
1363 std::io::stdout(),
1364 )
1365 .await
1366 {
1367 let error = startup_failure(
1368 StartupStage::Identity,
1369 "apply the owner identity startup policy",
1370 error,
1371 );
1372 return Err(close_writer_after_startup_failure(store, shutdown, task, error).await);
1373 }
1374
1375 let providers = ConfiguredLoginProviders::new(true, true)
1376 .expect("password and Nostr form a valid login-provider set");
1377 let security = match AdminSecurityState::new(
1378 host.admin_origin.clone(),
1379 store.auth.clone(),
1380 providers,
1381 AdminSessionPolicy::default(),
1382 Argon2idPolicy::v1(),
1383 )
1384 .await
1385 {
1386 Ok(security) => security,
1387 Err(error) => {
1388 let error = startup_failure(
1389 StartupStage::Identity,
1390 "initialize the admin authentication boundary",
1391 error,
1392 );
1393 return Err(close_writer_after_startup_failure(store, shutdown, task, error).await);
1394 }
1395 };
1396
1397 Ok(StartedDatabase {
1398 store,
1399 shutdown,
1400 task,
1401 security,
1402 })
1403}
1404
1405async fn close_started_database(
1406 database: StartedDatabase,
1407 process_error: ProcessError,
1408) -> ProcessError {
1409 close_writer_after_startup_failure(
1410 database.store,
1411 database.shutdown,
1412 database.task,
1413 process_error,
1414 )
1415 .await
1416}
1417
1418async fn close_writer_after_startup_failure(
1419 database: DatabaseStore,
1420 shutdown: CancellationToken,
1421 writer: JoinHandle<(CriticalTaskName, CriticalTaskCompletion)>,
1422 process_error: ProcessError,
1423) -> ProcessError {
1424 shutdown.cancel();
1425 drop(database);
1426 if let Some(error) = drained_task_failure(writer.await) {
1427 tracing::error!(error = %error, "database writer failed during startup cleanup");
1428 }
1429 process_error
1430}
1431
1432impl ApplicationRuntime {
1433 fn with_parts(
1434 readiness: Readiness,
1435 cancellation: CancellationToken,
1436 shutdown: ShutdownFuture,
1437 critical_tasks: Vec<CriticalTask>,
1438 ) -> Self {
1439 let mut supervisor = JoinSet::new();
1440 for task in critical_tasks {
1441 let span = task_span(task.name);
1442 supervisor.spawn(task.instrument(span));
1443 }
1444
1445 Self {
1446 readiness,
1447 cancellation,
1448 database_shutdown: CancellationToken::new(),
1449 shutdown,
1450 critical_tasks: supervisor,
1451 database_writer: None,
1452 }
1453 }
1454
1455 fn with_database_writer(
1456 readiness: Readiness,
1457 cancellation: CancellationToken,
1458 database_shutdown: CancellationToken,
1459 shutdown: ShutdownFuture,
1460 critical_tasks: Vec<CriticalTask>,
1461 database_writer: JoinHandle<(CriticalTaskName, CriticalTaskCompletion)>,
1462 ) -> Self {
1463 let mut runtime = Self::with_parts(readiness, cancellation, shutdown, critical_tasks);
1464 runtime.database_shutdown = database_shutdown;
1465 runtime.database_writer = Some(database_writer);
1466 runtime
1467 }
1468
1469 async fn run_until_stop(mut self) -> Result<(), ApplicationError> {
1470 self.readiness.mark_ready();
1471
1472 let failure = if self.critical_tasks.is_empty() && self.database_writer.is_none() {
1473 self.shutdown.await.err()
1474 } else {
1475 tokio::select! {
1476 biased;
1477 shutdown = &mut self.shutdown => shutdown.err(),
1478 completion = self.critical_tasks.join_next(), if !self.critical_tasks.is_empty() => {
1479 Some(unexpected_task_failure(completion))
1480 }
1481 completion = wait_for_database_writer(&mut self.database_writer) => {
1482 self.database_writer.take();
1483 Some(unexpected_task_failure(Some(completion)))
1484 }
1485 }
1486 };
1487
1488 self.readiness.mark_not_ready();
1489 self.cancellation.cancel();
1490
1491 let drain_failure = drain_critical_tasks(&mut self.critical_tasks).await;
1492 self.database_shutdown.cancel();
1493 let database_failure = drain_database_writer(&mut self.database_writer).await;
1494 match first_shutdown_failure([failure, drain_failure, database_failure]) {
1495 Some(error) => Err(error),
1496 None => Ok(()),
1497 }
1498 }
1499}
1500
1501fn spawn_critical_task(
1502 task: CriticalTask,
1503) -> JoinHandle<(CriticalTaskName, CriticalTaskCompletion)> {
1504 let span = task_span(task.name);
1505 tokio::spawn(task.instrument(span))
1506}
1507
1508fn first_shutdown_failure(failures: [Option<ApplicationError>; 3]) -> Option<ApplicationError> {
1509 let mut failures = failures.into_iter().flatten();
1510 let first = failures.next();
1511 for failure in failures {
1512 tracing::error!(error = %failure, "additional failure during shutdown");
1513 }
1514 first
1515}
1516
1517async fn wait_for_database_writer(
1518 writer: &mut Option<JoinHandle<(CriticalTaskName, CriticalTaskCompletion)>>,
1519) -> Result<(CriticalTaskName, CriticalTaskCompletion), JoinError> {
1520 match writer.as_mut() {
1521 Some(writer) => writer.await,
1522 None => std::future::pending().await,
1523 }
1524}
1525
1526async fn drain_database_writer(
1527 writer: &mut Option<JoinHandle<(CriticalTaskName, CriticalTaskCompletion)>>,
1528) -> Option<ApplicationError> {
1529 let writer = writer.take()?;
1530 drained_task_failure(writer.await)
1531}
1532
1533async fn drain_critical_tasks(
1534 critical_tasks: &mut JoinSet<(CriticalTaskName, CriticalTaskCompletion)>,
1535) -> Option<ApplicationError> {
1536 let mut first_failure = None;
1537
1538 while let Some(completion) = critical_tasks.join_next().await {
1539 if let Some(failure) = drained_task_failure(completion) {
1540 if first_failure.is_none() {
1541 first_failure = Some(failure);
1542 } else {
1543 tracing::error!(error = %failure, "additional critical task failure during shutdown");
1544 }
1545 }
1546 }
1547
1548 first_failure
1549}
1550
1551fn unexpected_task_failure(
1552 completion: Option<Result<(CriticalTaskName, CriticalTaskCompletion), JoinError>>,
1553) -> ApplicationError {
1554 match completion {
1555 Some(Ok((task, completion))) => task_completion_failure(task, completion)
1556 .unwrap_or(ApplicationError::CriticalTaskExited { task }),
1557 Some(Err(source)) => ApplicationError::TaskSupervisor { source },
1558 None => ApplicationError::TaskSupervisorEmpty,
1559 }
1560}
1561
1562fn drained_task_failure(
1563 completion: Result<(CriticalTaskName, CriticalTaskCompletion), JoinError>,
1564) -> Option<ApplicationError> {
1565 match completion {
1566 Ok((task, completion)) => task_completion_failure(task, completion),
1567 Err(source) => Some(ApplicationError::TaskSupervisor { source }),
1568 }
1569}
1570
1571fn task_completion_failure(
1572 task: CriticalTaskName,
1573 completion: CriticalTaskCompletion,
1574) -> Option<ApplicationError> {
1575 match completion {
1576 CriticalTaskCompletion::Returned(Ok(())) => None,
1577 CriticalTaskCompletion::Returned(Err(source)) => {
1578 Some(ApplicationError::CriticalTaskFailed { task, source })
1579 }
1580 CriticalTaskCompletion::Panicked(message) => {
1581 Some(ApplicationError::CriticalTaskPanicked { task, message })
1582 }
1583 }
1584}
1585
1586struct CriticalTask {
1587 name: CriticalTaskName,
1588 future: CriticalTaskFuture,
1589}
1590
1591impl CriticalTask {
1592 fn new<Future, Error>(name: CriticalTaskName, future: Future) -> Self
1593 where
1594 Future: std::future::Future<Output = Result<(), Error>> + Send + 'static,
1595 Error: Into<CriticalTaskFailure>,
1596 {
1597 Self {
1598 name,
1599 future: Box::pin(async move { future.await.map_err(Into::into) }),
1600 }
1601 }
1602}
1603
1604enum CriticalTaskCompletion {
1605 Returned(CriticalTaskResult),
1606 Panicked(Box<str>),
1607}
1608
1609impl Future for CriticalTask {
1610 type Output = (CriticalTaskName, CriticalTaskCompletion);
1611
1612 fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Self::Output> {
1613 let task = self.get_mut();
1614 match catch_unwind(AssertUnwindSafe(|| task.future.as_mut().poll(context))) {
1615 Ok(Poll::Pending) => Poll::Pending,
1616 Ok(Poll::Ready(result)) => {
1617 Poll::Ready((task.name, CriticalTaskCompletion::Returned(result)))
1618 }
1619 Err(payload) => Poll::Ready((
1620 task.name,
1621 CriticalTaskCompletion::Panicked(panic_message(payload.as_ref())),
1622 )),
1623 }
1624 }
1625}
1626
1627fn panic_message(payload: &(dyn std::any::Any + Send)) -> Box<str> {
1628 if let Some(message) = payload.downcast_ref::<String>() {
1629 return message.clone().into_boxed_str();
1630 }
1631 if let Some(message) = payload.downcast_ref::<&'static str>() {
1632 return Box::from(*message);
1633 }
1634 Box::from("non-string panic payload")
1635}
1636
1637#[cfg(unix)]
1638fn install_termination_signal() -> Result<ShutdownFuture, ApplicationError> {
1639 use tokio::signal::unix::{SignalKind, signal};
1640
1641 let interrupt =
1642 signal(SignalKind::interrupt()).map_err(|source| ApplicationError::SignalRegistration {
1643 signal: ShutdownSignal::Interrupt,
1644 source,
1645 })?;
1646 let terminate =
1647 signal(SignalKind::terminate()).map_err(|source| ApplicationError::SignalRegistration {
1648 signal: ShutdownSignal::Terminate,
1649 source,
1650 })?;
1651
1652 Ok(Box::pin(wait_for_unix_termination(interrupt, terminate)))
1653}
1654
1655#[cfg(unix)]
1656async fn wait_for_unix_termination(
1657 mut interrupt: tokio::signal::unix::Signal,
1658 mut terminate: tokio::signal::unix::Signal,
1659) -> Result<(), ApplicationError> {
1660 let closed_signal = tokio::select! {
1661 received = interrupt.recv() => {
1662 if received.is_some() {
1663 return Ok(());
1664 }
1665 ShutdownSignal::Interrupt
1666 }
1667 received = terminate.recv() => {
1668 if received.is_some() {
1669 return Ok(());
1670 }
1671 ShutdownSignal::Terminate
1672 }
1673 };
1674
1675 Err(ApplicationError::SignalStreamClosed {
1676 signal: closed_signal,
1677 })
1678}
1679
1680#[cfg(windows)]
1681fn install_termination_signal() -> Result<ShutdownFuture, ApplicationError> {
1682 let interrupt = tokio::signal::windows::ctrl_c().map_err(|source| {
1683 ApplicationError::SignalRegistration {
1684 signal: ShutdownSignal::Interrupt,
1685 source,
1686 }
1687 })?;
1688
1689 Ok(Box::pin(wait_for_windows_termination(interrupt)))
1690}
1691
1692#[cfg(windows)]
1693async fn wait_for_windows_termination(
1694 mut interrupt: tokio::signal::windows::CtrlC,
1695) -> Result<(), ApplicationError> {
1696 if interrupt.recv().await.is_some() {
1697 Ok(())
1698 } else {
1699 Err(ApplicationError::SignalStreamClosed {
1700 signal: ShutdownSignal::Interrupt,
1701 })
1702 }
1703}
1704
1705#[cfg(not(any(unix, windows)))]
1706fn install_termination_signal() -> Result<ShutdownFuture, ApplicationError> {
1707 Err(ApplicationError::SignalPlatformUnsupported)
1708}
1709
1710#[cfg(test)]
1711mod tests {
1712 use std::{
1713 cell::Cell,
1714 fs,
1715 path::PathBuf,
1716 sync::{
1717 Arc,
1718 atomic::{AtomicBool, AtomicUsize, Ordering},
1719 },
1720 };
1721
1722 use axum::http::{Method, StatusCode};
1723 use k256::schnorr::SigningKey;
1724 use markdown_compiler::DefaultPostTipPolicy;
1725 use serde::{Serialize, de::DeserializeOwned};
1726 use tokio::{
1727 io::{AsyncReadExt as _, AsyncWriteExt as _},
1728 net::{TcpSocket, TcpStream},
1729 sync::oneshot,
1730 };
1731
1732 use maincopy_shared::auth::{
1733 AdminAuditEventId, AdminScope, AgentCredentialId, InstanceId, UserId,
1734 };
1735
1736 use crate::{
1737 admin::test_support::{ADMIN_AUTHORITY, ADMIN_ORIGIN, agent_authorization},
1738 cli::BootstrapCredential,
1739 config::{ConfigurationValidationCode, IdentityStartupBootstrap},
1740 domain::{
1741 auth::{
1742 NostrPublicKey,
1743 store::{
1744 AdminMutationKey, AuditPrincipalReference, BootstrapIdentity,
1745 MutationAuditContext, NewHumanCredential, RegisterAgentCredential,
1746 },
1747 },
1748 mail::runtime::{MailReviewAccessError, MailStartupError},
1749 },
1750 render::render_bound_post_preview,
1751 };
1752
1753 use super::*;
1754
1755 const VALID_PUBLICATION: &str = "[site]\n\
1756title = \"Pinned startup source\"\n\
1757base_url = \"https://startup.example.test\"\n\
1758description = \"Startup configuration test.\"\n\
1759[author]\n\
1760name = \"Startup Tester\"\n";
1761 const DURABLE_POST_ID: &str = "11111111-1111-4111-8111-111111111111";
1762 const DURABLE_PUBLICATION_ID: &str = "aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa";
1763 const LIVE_RELOAD_TEST_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(15);
1764 const TEST_OWNER_NOSTR_PUBLIC_KEY: &str =
1765 "63fe6318dc58583cfe16810f86dd09e18bfd76aabc24a0081ce2856f330504ed";
1766 const TEST_AGENT_SECRET: [u8; 32] = [3_u8; 32];
1767 const DURABLE_POST: &str = "+++\n\
1768id = \"11111111-1111-4111-8111-111111111111\"\n\
1769title = \"Durable publication\"\n\
1770slug = \"durable-publication\"\n\
1771authored_at = 2026-08-29T12:00:00Z\n\
1772description = \"A publication restored from SQLite.\"\n\
1773+++\n\
1774# Durable publication\n\n\
1775Durable article body.\n";
1776
1777 fn reserve_loopback_port() -> TcpSocket {
1778 let reservation = TcpSocket::new_v4().unwrap();
1781 reservation.set_reuseaddr(true).unwrap();
1782 reservation.bind("127.0.0.1:0".parse().unwrap()).unwrap();
1783 reservation
1784 }
1785
1786 fn startup_host_source(extra: &str, public_bind: &str) -> String {
1787 startup_host_source_with_listeners(extra, public_bind, "127.0.0.1:0", "127.0.0.1:0")
1788 }
1789
1790 fn startup_host_source_with_listeners(
1791 extra: &str,
1792 public_bind: &str,
1793 admin_bind: &str,
1794 metrics_bind: &str,
1795 ) -> String {
1796 format!(
1797 "[paths]\n\
1798 content_root = \"content\"\n\
1799 state_root = \"state\"\n\
1800 runtime_root = \"run\"\n\
1801 [public]\n\
1802 bind = \"{public_bind}\"\n\
1803 {extra}\n\
1804 [metrics]\n\
1805 bind = \"{metrics_bind}\"\n\
1806 [admin]\n\
1807 bind = \"{admin_bind}\"\n\
1808 origin = \"https://admin.example.test\"\n"
1809 )
1810 }
1811
1812 fn startup_fixture(
1813 host_source: &str,
1814 publication_source: &str,
1815 ) -> (tempfile::TempDir, PathBuf, PathBuf) {
1816 let root = tempfile::tempdir().unwrap();
1817 let content_root = root.path().join("content");
1818 fs::create_dir(&content_root).unwrap();
1819 let publication_path = content_root.join("publication.toml");
1820 fs::write(&publication_path, publication_source).unwrap();
1821 let config_path = root.path().join("maincopy.toml");
1822 fs::write(
1823 &config_path,
1824 startup_host_source(host_source, "127.0.0.1:0"),
1825 )
1826 .unwrap();
1827 (root, config_path, publication_path)
1828 }
1829
1830 fn write_durable_post(content_root: &Path) {
1831 fs::create_dir(content_root.join("posts")).unwrap();
1832 fs::write(
1833 content_root.join("posts/durable-publication.md"),
1834 DURABLE_POST,
1835 )
1836 .unwrap();
1837 }
1838
1839 struct AdminTestClient {
1840 address: std::net::SocketAddr,
1841 client: reqwest::Client,
1842 signing_key: SigningKey,
1843 }
1844
1845 fn admin_client(application: &Application) -> AdminTestClient {
1846 AdminTestClient {
1847 address: application.admin_addr,
1848 client: reqwest::Client::builder()
1849 .no_proxy()
1850 .timeout(std::time::Duration::from_secs(5))
1851 .build()
1852 .unwrap(),
1853 signing_key: SigningKey::from_bytes(&TEST_AGENT_SECRET).unwrap(),
1854 }
1855 }
1856
1857 async fn admin_get(admin: &AdminTestClient, path: &str) -> reqwest::Response {
1858 let idempotency_key = uuid::Uuid::new_v4().to_string();
1859 let authorization = agent_authorization(
1860 &admin.signing_key,
1861 &Method::GET,
1862 path,
1863 &[],
1864 Some(&idempotency_key),
1865 );
1866 admin
1867 .client
1868 .get(format!("http://{}{path}", admin.address))
1869 .header(reqwest::header::HOST, ADMIN_AUTHORITY)
1870 .header("origin", ADMIN_ORIGIN)
1871 .header(reqwest::header::AUTHORIZATION, authorization)
1872 .header(
1873 maincopy_shared::publication::IDEMPOTENCY_KEY_HEADER,
1874 idempotency_key,
1875 )
1876 .send()
1877 .await
1878 .unwrap()
1879 }
1880
1881 async fn admin_post_json(
1882 admin: &AdminTestClient,
1883 path: &str,
1884 idempotency_key: &str,
1885 value: &impl Serialize,
1886 ) -> reqwest::Response {
1887 admin_post_body(
1888 admin,
1889 path,
1890 idempotency_key,
1891 None,
1892 serde_json::to_vec(value).unwrap(),
1893 )
1894 .await
1895 }
1896
1897 async fn admin_post_body(
1898 admin: &AdminTestClient,
1899 path: &str,
1900 idempotency_key: &str,
1901 request_id: Option<&str>,
1902 body: impl AsRef<[u8]>,
1903 ) -> reqwest::Response {
1904 let body = body.as_ref().to_vec();
1905 let authorization = agent_authorization(
1906 &admin.signing_key,
1907 &Method::POST,
1908 path,
1909 &body,
1910 Some(idempotency_key),
1911 );
1912 let mut request = admin
1913 .client
1914 .post(format!("http://{}{path}", admin.address))
1915 .header(reqwest::header::HOST, ADMIN_AUTHORITY)
1916 .header("origin", ADMIN_ORIGIN)
1917 .header(reqwest::header::AUTHORIZATION, authorization)
1918 .header(
1919 maincopy_shared::publication::IDEMPOTENCY_KEY_HEADER,
1920 idempotency_key,
1921 )
1922 .header(reqwest::header::CONTENT_TYPE, "application/json");
1923 if let Some(request_id) = request_id {
1924 request = request.header("x-request-id", request_id);
1925 }
1926 request.body(body).send().await.unwrap()
1927 }
1928
1929 async fn admin_json<ResponseType: DeserializeOwned>(
1930 response: reqwest::Response,
1931 ) -> ResponseType {
1932 let bytes = response.bytes().await.unwrap();
1933 serde_json::from_slice(&bytes).unwrap()
1934 }
1935
1936 async fn admin_text(response: reqwest::Response) -> String {
1937 response.text().await.unwrap()
1938 }
1939
1940 async fn admin_preview_digest(
1941 admin: &AdminTestClient,
1942 post_id: uuid::Uuid,
1943 ) -> maincopy_shared::publication::PreviewDigest {
1944 use maincopy_shared::{
1945 posts::POSTS_PATH,
1946 publication::{PREVIEW_DIGEST_HEADER, PreviewDigest},
1947 };
1948
1949 let response = admin_get(admin, &format!("{POSTS_PATH}/{post_id}/preview")).await;
1950 assert_eq!(response.status(), StatusCode::OK);
1951 let encoded = response
1952 .headers()
1953 .get(PREVIEW_DIGEST_HEADER)
1954 .unwrap()
1955 .to_str()
1956 .unwrap();
1957 PreviewDigest::parse(encoded).unwrap()
1958 }
1959
1960 async fn stop_built_application(application: Application) {
1961 let Application {
1962 _process_lock: process_lock,
1963 _database: database,
1964 publication_coordinator,
1965 mut runtime,
1966 public_addr: _,
1967 admin_addr: _,
1968 } = application;
1969 runtime.cancellation.cancel();
1970 assert!(
1971 drain_critical_tasks(&mut runtime.critical_tasks)
1972 .await
1973 .is_none()
1974 );
1975 runtime.database_shutdown.cancel();
1976 assert!(
1977 drain_database_writer(&mut runtime.database_writer)
1978 .await
1979 .is_none()
1980 );
1981 drop(publication_coordinator);
1982 drop(database);
1983 drop(process_lock);
1984 tokio::task::yield_now().await;
1985 }
1986
1987 async fn bootstrap_test_identity(startup: &StartupConfiguration) {
1988 let host = startup._host.view();
1989 let database = database::bootstrap(host.database).await.unwrap();
1990 let (store, writer) = database.into_store(host.database.writer_queue_capacity.get());
1991 let shutdown = CancellationToken::new();
1992 let writer_shutdown = shutdown.clone();
1993 let writer = tokio::spawn(async move {
1994 writer.run(writer_shutdown).await.unwrap();
1995 });
1996 if store
1997 .auth
1998 .identity_state()
1999 .await
2000 .unwrap()
2001 .bootstrap_required
2002 {
2003 let owner_user_id = UserId::from_uuid(uuid::Uuid::new_v4());
2004 store
2005 .auth
2006 .bootstrap_identity(BootstrapIdentity {
2007 instance_id: InstanceId::from_uuid(uuid::Uuid::new_v4()),
2008 owner_user_id,
2009 credential: NewHumanCredential::Nostr {
2010 public_key: NostrPublicKey::parse(TEST_OWNER_NOSTR_PUBLIC_KEY).unwrap(),
2011 },
2012 configured_providers: ConfiguredLoginProviders::new(true, true).unwrap(),
2013 occurred_at: OffsetDateTime::now_utc(),
2014 audit_event_id: AdminAuditEventId::from_uuid(uuid::Uuid::new_v4()),
2015 })
2016 .await
2017 .unwrap();
2018 let signing_key = SigningKey::from_bytes(&TEST_AGENT_SECRET).unwrap();
2019 store
2020 .auth
2021 .register_agent_credential(RegisterAgentCredential {
2022 credential_id: AgentCredentialId::from_uuid(uuid::Uuid::new_v4()),
2023 owner_user_id,
2024 issuer_user_id: owner_user_id,
2025 public_key: NostrPublicKey::from_bytes(
2026 signing_key.verifying_key().to_bytes().into(),
2027 )
2028 .unwrap(),
2029 label: "startup integration agent".into(),
2030 scopes: AdminScope::PUBLISHER.into_iter().collect(),
2031 created_at: OffsetDateTime::now_utc(),
2032 expires_at: None,
2033 audit: MutationAuditContext {
2034 audit_event_id: AdminAuditEventId::from_uuid(uuid::Uuid::new_v4()),
2035 principal: AuditPrincipalReference::Offline {
2036 user_id: Some(owner_user_id),
2037 },
2038 request_id: None,
2039 idempotency_key: AdminMutationKey(uuid::Uuid::new_v4()),
2040 },
2041 })
2042 .await
2043 .unwrap();
2044 }
2045 drop(store);
2046 shutdown.cancel();
2047 writer.await.unwrap();
2048 }
2049
2050 async fn build_test_application(
2051 startup: StartupConfiguration,
2052 ) -> Result<Application, ProcessError> {
2053 bootstrap_test_identity(&startup).await;
2054 Application::build(startup).await
2055 }
2056
2057 #[test]
2058 #[cfg(target_os = "linux")]
2059 fn content_discovery_runs_once_receives_effective_limits_and_validation_reuses_owned_bytes() {
2060 let host_source = "[content]\n\
2061publication_file_bytes = 512\n\
2062post_file_bytes = 1536\n\
2063asset_file_bytes = 2560\n\
2064total_tree_bytes = 3584\n\
2065entries = 90\n\
2066depth = 7\n\
2067path_bytes = 400\n";
2068 let (_root, arguments, publication_path) = startup_fixture(host_source, VALID_PUBLICATION);
2069 let discovery_calls = Cell::new(0);
2070 let observed_limits = Cell::new(None);
2071
2072 let startup =
2073 StartupConfiguration::load_with_discovery(arguments, |content_root, limits| {
2074 discovery_calls.set(discovery_calls.get() + 1);
2075 observed_limits.set(Some(limits));
2076 let tree = discover_content_tree(content_root, limits)?;
2077 fs::write(&publication_path, "not the discovered publication").unwrap();
2078 Ok(tree)
2079 })
2080 .unwrap();
2081
2082 assert_eq!(discovery_calls.get(), 1);
2083 let limits = observed_limits.get().unwrap();
2084 assert_eq!(limits.publication_file_bytes.get(), 512);
2085 assert_eq!(limits.post_file_bytes.get(), 1536);
2086 assert_eq!(limits.asset_file_bytes.get(), 2560);
2087 assert_eq!(limits.total_tree_bytes.get(), 3584);
2088 assert_eq!(limits.entries.get(), 90);
2089 assert_eq!(limits.depth.get(), 7);
2090 assert_eq!(limits.path_bytes.get(), 400);
2091 assert!(
2092 startup
2093 ._content_tree
2094 .publication
2095 .source
2096 .contains("Pinned startup source")
2097 );
2098 assert_eq!(
2099 startup.content.view().publication.site.title.as_str(),
2100 "Pinned startup source"
2101 );
2102 }
2103
2104 #[test]
2105 fn host_configuration_failure_prevents_content_discovery() {
2106 let root = tempfile::tempdir().unwrap();
2107 let arguments = root.path().join("missing-maincopy.toml");
2108 let discovery_calls = Cell::new(0);
2109
2110 let result = StartupConfiguration::load_with_discovery(arguments, |_, _| {
2111 discovery_calls.set(discovery_calls.get() + 1);
2112 unreachable!("content discovery must follow host configuration")
2113 });
2114
2115 assert!(matches!(result, Err(ProcessError::Configuration(_))));
2116 assert_eq!(discovery_calls.get(), 0);
2117 }
2118
2119 #[test]
2120 #[cfg(target_os = "linux")]
2121 fn content_validation_failure_has_stable_exit() {
2122 let (_root, arguments, _) = startup_fixture("", "unknown = true\n");
2123 let discovery_calls = Cell::new(0);
2124
2125 let result = StartupConfiguration::load_with_discovery(arguments, |root, limits| {
2126 discovery_calls.set(discovery_calls.get() + 1);
2127 discover_content_tree(root, limits)
2128 });
2129
2130 let Err(error) = result else {
2131 panic!("invalid content must fail startup");
2132 };
2133 assert!(matches!(error, ProcessError::Validation(_)));
2134 assert_eq!(error.exit(), ProcessExit::Validation);
2135 assert_eq!(discovery_calls.get(), 1);
2136 }
2137
2138 #[test]
2139 #[cfg(target_os = "linux")]
2140 fn authored_tips_do_not_require_a_payment_provider() {
2141 let publication = format!(
2142 "{VALID_PUBLICATION}[tips]\n\
2143 enabled = true\n"
2144 );
2145 let (_root, arguments, _) = startup_fixture("", &publication);
2146
2147 let startup =
2148 StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
2149
2150 assert_eq!(
2151 startup.content.view().publication.tips,
2152 DefaultPostTipPolicy::Enabled
2153 );
2154 }
2155
2156 #[test]
2157 #[cfg(target_os = "linux")]
2158 fn removed_provider_configuration_is_rejected_without_opening_credentials() {
2159 let host = "[lightning]\n\
2160provider = \"lexe\"\n\
2161network = \"signet\"\n\
2162credentials = { source = \"file\", path = \"must-not-open.json\" }\n";
2163 let (root, arguments, _) = startup_fixture(host, VALID_PUBLICATION);
2164 let credential_path = root.path().join("must-not-open.json");
2165
2166 let Err(error) =
2167 StartupConfiguration::load_with_discovery(arguments, discover_content_tree)
2168 else {
2169 panic!("removed provider configuration must fail host validation");
2170 };
2171
2172 let ProcessError::Configuration(errors) = error else {
2173 panic!("removed provider configuration must fail host validation");
2174 };
2175 assert_eq!(
2176 errors.diagnostics()[0].code,
2177 ConfigurationValidationCode::HostTomlInvalid
2178 );
2179 assert!(!credential_path.exists());
2180 }
2181
2182 #[tokio::test]
2183 async fn production_identity_policy_refuses_credential_output_and_accepts_offline_nostr_bootstrap()
2184 {
2185 let (_root, config_path, _) = startup_fixture(
2186 "[identity]\nstartup_bootstrap = \"require_existing\"\n",
2187 VALID_PUBLICATION,
2188 );
2189 let host = HostConfigurationLoader::from_process_working_directory()
2190 .unwrap()
2191 .load(&config_path)
2192 .unwrap();
2193 let database = database::bootstrap(host.view().database).await.unwrap();
2194 let (store, writer) = database.into_store(host.view().database.writer_queue_capacity.get());
2195 let cancellation = CancellationToken::new();
2196 let running = tokio::spawn(writer.run(cancellation.clone()));
2197 let mut output = Vec::new();
2198 let failure = identity_bootstrap::initialize_startup_identity(
2199 &store,
2200 IdentityStartupBootstrap::RequireExisting,
2201 &mut output,
2202 )
2203 .await
2204 .unwrap_err();
2205 assert!(matches!(
2206 failure,
2207 identity_bootstrap::GeneratedOwnerBootstrapError::ExistingIdentityRequired
2208 ));
2209 assert!(output.is_empty());
2210 assert!(
2211 store
2212 .auth
2213 .identity_state()
2214 .await
2215 .unwrap()
2216 .bootstrap_required
2217 );
2218 cancellation.cancel();
2219 running.await.unwrap().unwrap();
2220 drop(store);
2221 let startup =
2222 StartupConfiguration::load_with_discovery(config_path.clone(), discover_content_tree)
2223 .unwrap();
2224 let failure = match Application::build(startup).await {
2225 Ok(application) => {
2226 stop_built_application(application).await;
2227 panic!("production startup must refuse an uninitialized identity");
2228 }
2229 Err(error) => error,
2230 };
2231 assert!(matches!(
2232 failure,
2233 ProcessError::Application(ApplicationError::Startup {
2234 stage: StartupStage::Identity,
2235 operation: "apply the owner identity startup policy",
2236 ..
2237 })
2238 ));
2239 let public_key = NostrPublicKey::parse(
2240 "f9308a019258c31049344f85f89d5229b531c845836f99b08601f113bce036f9",
2241 )
2242 .unwrap();
2243 identity_bootstrap::bootstrap_owner(
2244 config_path.clone(),
2245 BootstrapCredential::Nostr { public_key },
2246 )
2247 .await
2248 .unwrap();
2249 let startup =
2250 StartupConfiguration::load_with_discovery(config_path, discover_content_tree).unwrap();
2251 let application = Application::build(startup).await.unwrap();
2252 let owner = application
2253 ._database
2254 .auth
2255 .users_page(None, 2)
2256 .await
2257 .unwrap();
2258 assert_eq!(owner.items.len(), 1);
2259 assert!(owner.items[0].has_nostr);
2260 assert!(!owner.items[0].has_password);
2261 let mut output = Vec::new();
2262 assert!(
2263 !identity_bootstrap::initialize_startup_identity(
2264 &application._database,
2265 IdentityStartupBootstrap::RequireExisting,
2266 &mut output
2267 )
2268 .await
2269 .unwrap()
2270 );
2271 assert!(output.is_empty());
2272 stop_built_application(application).await;
2273 }
2274
2275 #[tokio::test]
2276 #[cfg(target_os = "linux")]
2277 async fn unbootstrapped_identity_generates_owner_and_starts_the_application() {
2278 let (root, _, _) = startup_fixture("", VALID_PUBLICATION);
2279 let config_path = root.path().join("maincopy.toml");
2280 let reservation = reserve_loopback_port();
2281 let public_addr = reservation.local_addr().unwrap();
2282 fs::write(
2283 &config_path,
2284 startup_host_source("", &public_addr.to_string()),
2285 )
2286 .unwrap();
2287 let arguments = config_path;
2288 let startup =
2289 StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
2290 let application = Application::build(startup).await.unwrap();
2291
2292 let identity = application._database.auth.identity_state().await.unwrap();
2293 assert!(!identity.bootstrap_required);
2294 assert!(identity.instance.is_some());
2295 let users = application
2296 ._database
2297 .auth
2298 .users_page(None, 2)
2299 .await
2300 .unwrap();
2301 assert_eq!(users.items.len(), 1);
2302 let owner = &users.items[0];
2303 assert!(owner.has_password);
2304 assert!(!owner.has_nostr);
2305 assert!(
2306 owner
2307 .roles
2308 .contains(&maincopy_shared::auth::UserRole::Owner)
2309 );
2310 let credentials = application
2311 ._database
2312 .auth
2313 .user_credentials(owner.user_id)
2314 .await
2315 .unwrap()
2316 .unwrap();
2317 assert!(matches!(
2318 credentials.as_slice(),
2319 [crate::domain::auth::store::StoredHumanCredential::Password {
2320 username,
2321 ..
2322 }] if username.as_str() == "owner"
2323 ));
2324 assert!(
2325 !identity_bootstrap::initialize_startup_identity(
2326 &application._database,
2327 IdentityStartupBootstrap::GenerateOwner,
2328 std::io::sink()
2329 )
2330 .await
2331 .unwrap()
2332 );
2333
2334 stop_built_application(application).await;
2335 drop(reservation);
2336 }
2337
2338 #[tokio::test]
2339 #[cfg(target_os = "linux")]
2340 async fn application_build_serves_public_site_and_protected_admin_backend() {
2341 let (root, arguments, _) = startup_fixture("", VALID_PUBLICATION);
2342 let startup =
2343 StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
2344
2345 let application = build_test_application(startup).await.unwrap();
2346
2347 assert!(matches!(
2348 StartupHostConfiguration::load(root.path().join("maincopy.toml")),
2349 Err(ProcessError::AlreadyRunning)
2350 ));
2351 assert_eq!(application.runtime.critical_tasks.len(), 8);
2353 assert!(application.runtime.database_writer.is_some());
2354 let client = reqwest::Client::builder()
2355 .no_proxy()
2356 .timeout(std::time::Duration::from_secs(2))
2357 .build()
2358 .unwrap();
2359 let response = client
2360 .get(format!("http://{}/", application.public_addr))
2361 .send()
2362 .await
2363 .unwrap();
2364 assert_eq!(response.status(), reqwest::StatusCode::OK);
2365 let html = response.text().await.unwrap();
2366 assert!(html.contains("Pinned startup source"));
2367 assert!(html.contains("No posts have been published yet."));
2368
2369 let protected = client
2370 .get(format!(
2371 "http://{}/api/admin/v1/capabilities",
2372 application.admin_addr
2373 ))
2374 .header(reqwest::header::HOST, "admin.example.test")
2375 .send()
2376 .await
2377 .unwrap();
2378 assert_eq!(protected.status(), reqwest::StatusCode::UNAUTHORIZED);
2379 assert_eq!(
2380 protected.headers()[reqwest::header::CACHE_CONTROL],
2381 "private, no-store"
2382 );
2383
2384 stop_built_application(application).await;
2385 assert!(StartupHostConfiguration::load(root.path().join("maincopy.toml")).is_ok());
2386 }
2387
2388 #[tokio::test]
2389 #[cfg(target_os = "linux")]
2390 async fn admin_publication_route_activates_the_public_site_and_replays_success() {
2391 use maincopy_shared::publication::{
2392 PUBLICATIONS_PATH, PublishNowRequest, PublishNowResponse,
2393 };
2394
2395 let (root, arguments, _) = startup_fixture("", VALID_PUBLICATION);
2396 let content_root = root.path().join("content");
2397 write_durable_post(&content_root);
2398 let startup =
2399 StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
2400 let application = build_test_application(startup).await.unwrap();
2401 let admin = admin_client(&application);
2402 let post_id = uuid::Uuid::parse_str(DURABLE_POST_ID).unwrap();
2403 let request = PublishNowRequest {
2404 post_id,
2405 preview_digest: admin_preview_digest(&admin, post_id).await,
2406 expected_revision: None,
2407 scheduled_for: None,
2408 };
2409 let creation_key = "bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb";
2410
2411 let response = admin_post_json(&admin, PUBLICATIONS_PATH, creation_key, &request).await;
2412 assert_eq!(response.status(), StatusCode::OK);
2413 let published: PublishNowResponse = admin_json(response).await;
2414 assert_eq!(published.post_id, request.post_id);
2415 assert!(published.revision.starts_with("post-b3-v1-"));
2416 assert!(published.site_digest.starts_with("site-b3-v1-"));
2417 assert_eq!(published.site_version, 2);
2418
2419 let replay = admin_post_body(
2420 &admin,
2421 PUBLICATIONS_PATH,
2422 creation_key,
2423 None,
2424 serde_json::to_string_pretty(&request).unwrap(),
2425 )
2426 .await;
2427 assert_eq!(replay.status(), StatusCode::OK);
2428 assert_eq!(admin_json::<PublishNowResponse>(replay).await, published);
2429
2430 let request_id = "67e55044-10b1-426f-9247-bb680e5fe0c8";
2431 let malformed_key = uuid::Uuid::new_v4().to_string();
2432 let malformed = admin_post_body(
2433 &admin,
2434 PUBLICATIONS_PATH,
2435 &malformed_key,
2436 Some(request_id),
2437 "{",
2438 )
2439 .await;
2440 assert_eq!(malformed.status(), StatusCode::BAD_REQUEST);
2441 let error: serde_json::Value = admin_json(malformed).await;
2442 assert_eq!(error["error"]["code"], "invalid_request_body");
2443 assert_eq!(error["error"]["request_id"], request_id);
2444
2445 let oversized_key = uuid::Uuid::new_v4().to_string();
2446 let oversized = admin_post_body(
2447 &admin,
2448 PUBLICATIONS_PATH,
2449 &oversized_key,
2450 Some(request_id),
2451 " ".repeat(4 * 1024 + 1),
2452 )
2453 .await;
2454 assert_eq!(oversized.status(), StatusCode::PAYLOAD_TOO_LARGE);
2455 assert_eq!(oversized.headers()["x-request-id"], request_id);
2456 let error: serde_json::Value = admin_json(oversized).await;
2457 assert_eq!(error["error"]["code"], "request_body_too_large");
2458 assert_eq!(error["error"]["request_id"], request_id);
2459 assert_eq!(
2460 admin_get(&admin, "/api/admin/v1/capabilities")
2461 .await
2462 .status(),
2463 StatusCode::OK
2464 );
2465
2466 let public = reqwest::get(format!(
2467 "http://{}/posts/durable-publication",
2468 application.public_addr
2469 ))
2470 .await
2471 .unwrap();
2472 assert_eq!(public.status(), reqwest::StatusCode::OK);
2473 assert!(
2474 public
2475 .text()
2476 .await
2477 .unwrap()
2478 .contains("Durable article body.")
2479 );
2480
2481 stop_built_application(application).await;
2482 }
2483
2484 #[tokio::test]
2485 #[cfg(target_os = "linux")]
2486 async fn edited_markdown_updates_private_preview_until_explicit_publication_approval() {
2487 use maincopy_shared::{
2488 posts::{ListPostsResponse, POSTS_PATH, PostPublicationState},
2489 publication::{PUBLICATIONS_PATH, PublishNowRequest, PublishNowResponse},
2490 };
2491
2492 let (root, arguments, _) = startup_fixture("", VALID_PUBLICATION);
2493 let content_root = root.path().join("content");
2494 write_durable_post(&content_root);
2495 let post_path = content_root.join("posts/durable-publication.md");
2496 let startup =
2497 StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
2498 let application = build_test_application(startup).await.unwrap();
2499 let admin = admin_client(&application);
2500 let response = admin_post_json(
2501 &admin,
2502 PUBLICATIONS_PATH,
2503 "bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb",
2504 &PublishNowRequest {
2505 post_id: uuid::Uuid::parse_str(DURABLE_POST_ID).unwrap(),
2506 preview_digest: admin_preview_digest(
2507 &admin,
2508 uuid::Uuid::parse_str(DURABLE_POST_ID).unwrap(),
2509 )
2510 .await,
2511 expected_revision: None,
2512 scheduled_for: None,
2513 },
2514 )
2515 .await;
2516 assert_eq!(response.status(), StatusCode::OK);
2517 let published: PublishNowResponse = admin_json(response).await;
2518
2519 let public_url = format!(
2520 "http://{}/posts/durable-publication",
2521 application.public_addr
2522 );
2523 let public = reqwest::Client::builder().no_proxy().build().unwrap();
2524 let initial = public.get(&public_url).send().await.unwrap();
2525 let initial_etag = initial.headers()[reqwest::header::ETAG].clone();
2526 assert!(
2527 initial
2528 .text()
2529 .await
2530 .unwrap()
2531 .contains("Durable article body.")
2532 );
2533
2534 let edited = DURABLE_POST.replace(
2535 "Durable article body.",
2536 "This edit appeared without restarting Maincopy.",
2537 );
2538 fs::write(&post_path, &edited).unwrap();
2539
2540 let edited_summary = tokio::time::timeout(LIVE_RELOAD_TEST_TIMEOUT, async {
2541 loop {
2542 let posts: ListPostsResponse =
2543 admin_json(admin_get(&admin, POSTS_PATH).await).await;
2544 if let Some(summary) = posts.posts.into_iter().find(|summary| {
2545 summary.post_id.to_string() == DURABLE_POST_ID
2546 && summary.publication_state == PostPublicationState::UnpublishedChange
2547 && summary.revision != published.revision
2548 }) {
2549 break summary;
2550 }
2551 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
2552 }
2553 })
2554 .await
2555 .expect("a stable Markdown edit must enter the private preview catalog");
2556 let preview_url = format!("{POSTS_PATH}/{}/preview", edited_summary.post_id);
2557 let preview = admin_get(&admin, &preview_url).await;
2558 assert_eq!(preview.status(), StatusCode::OK);
2559 let preview_body = admin_text(preview).await;
2560 assert!(preview_body.contains("This edit appeared without restarting Maincopy."));
2561 assert!(!preview_body.contains("Durable article body."));
2562
2563 let still_pinned = public.get(&public_url).send().await.unwrap();
2564 assert_eq!(still_pinned.headers()[reqwest::header::ETAG], initial_etag);
2565 let still_pinned_body = still_pinned.text().await.unwrap();
2566 assert!(still_pinned_body.contains("Durable article body."));
2567 assert!(!still_pinned_body.contains("This edit appeared without restarting Maincopy."));
2568
2569 fs::write(&post_path, "+++\ninvalid = true\n").unwrap();
2570 tokio::time::sleep(std::time::Duration::from_millis(1_500)).await;
2571 let posts: ListPostsResponse = admin_json(admin_get(&admin, POSTS_PATH).await).await;
2572 let after_invalid = posts
2573 .posts
2574 .into_iter()
2575 .find(|summary| summary.post_id == edited_summary.post_id)
2576 .unwrap();
2577 assert_eq!(after_invalid.revision, edited_summary.revision);
2578 assert_eq!(
2579 after_invalid.publication_state,
2580 PostPublicationState::UnpublishedChange
2581 );
2582 let preview = admin_get(&admin, &preview_url).await;
2583 assert_eq!(preview.status(), StatusCode::OK);
2584 assert!(
2585 admin_text(preview)
2586 .await
2587 .contains("This edit appeared without restarting Maincopy.")
2588 );
2589 let after_invalid_public = public.get(&public_url).send().await.unwrap();
2590 assert_eq!(
2591 after_invalid_public.headers()[reqwest::header::ETAG],
2592 initial_etag
2593 );
2594 let after_invalid_public_body = after_invalid_public.text().await.unwrap();
2595 assert!(after_invalid_public_body.contains("Durable article body."));
2596 assert!(
2597 !after_invalid_public_body.contains("This edit appeared without restarting Maincopy.")
2598 );
2599
2600 fs::write(&post_path, &edited).unwrap();
2601 let response = admin_post_json(
2602 &admin,
2603 PUBLICATIONS_PATH,
2604 "cccccccc-cccc-4ccc-8ccc-cccccccccccc",
2605 &PublishNowRequest {
2606 post_id: edited_summary.post_id,
2607 preview_digest: admin_preview_digest(&admin, edited_summary.post_id).await,
2608 expected_revision: Some(edited_summary.revision.clone()),
2609 scheduled_for: None,
2610 },
2611 )
2612 .await;
2613 assert_eq!(response.status(), StatusCode::OK);
2614 let approved: PublishNowResponse = admin_json(response).await;
2615 assert_eq!(approved.revision, edited_summary.revision);
2616
2617 let updated_public = public.get(&public_url).send().await.unwrap();
2618 assert_ne!(
2619 updated_public.headers()[reqwest::header::ETAG],
2620 initial_etag
2621 );
2622 let updated_body = updated_public.text().await.unwrap();
2623 assert!(updated_body.contains("This edit appeared without restarting Maincopy."));
2624 assert!(!updated_body.contains("Durable article body."));
2625
2626 stop_built_application(application).await;
2627 }
2628
2629 #[tokio::test]
2630 #[cfg(target_os = "linux")]
2631 async fn a_new_markdown_file_becomes_publishable_without_restarting() {
2632 use maincopy_shared::{
2633 posts::{ListPostsResponse, POSTS_PATH},
2634 publication::{PUBLICATIONS_PATH, PublishNowRequest},
2635 };
2636
2637 let (root, arguments, _) = startup_fixture("", VALID_PUBLICATION);
2638 let content_root = root.path().join("content");
2639 let startup =
2640 StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
2641 let application = build_test_application(startup).await.unwrap();
2642 let admin = admin_client(&application);
2643
2644 write_durable_post(&content_root);
2645 let summary = tokio::time::timeout(LIVE_RELOAD_TEST_TIMEOUT, async {
2646 loop {
2647 let posts: ListPostsResponse =
2648 admin_json(admin_get(&admin, POSTS_PATH).await).await;
2649 if let Some(post) = posts
2650 .posts
2651 .into_iter()
2652 .find(|post| post.post_id.to_string() == DURABLE_POST_ID)
2653 {
2654 break post;
2655 }
2656 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
2657 }
2658 })
2659 .await
2660 .expect("a stable new Markdown file must enter the live catalog");
2661
2662 let response = admin_post_json(
2663 &admin,
2664 PUBLICATIONS_PATH,
2665 "bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb",
2666 &PublishNowRequest {
2667 post_id: summary.post_id,
2668 preview_digest: admin_preview_digest(&admin, summary.post_id).await,
2669 expected_revision: Some(summary.revision),
2670 scheduled_for: None,
2671 },
2672 )
2673 .await;
2674 assert_eq!(response.status(), StatusCode::OK);
2675
2676 let body = reqwest::get(format!(
2677 "http://{}/posts/durable-publication",
2678 application.public_addr
2679 ))
2680 .await
2681 .unwrap()
2682 .text()
2683 .await
2684 .unwrap();
2685 assert!(body.contains("Durable article body."));
2686
2687 stop_built_application(application).await;
2688 }
2689
2690 #[tokio::test]
2691 #[cfg(target_os = "linux")]
2692 async fn restart_keeps_public_revision_and_aliases_pinned_until_stopped_edit_is_approved() {
2693 use maincopy_shared::{
2694 posts::{ListPostsResponse, POSTS_PATH, PostPublicationState},
2695 publication::{PUBLICATIONS_PATH, PublishNowRequest, PublishNowResponse},
2696 };
2697
2698 let (root, arguments, _) = startup_fixture("", VALID_PUBLICATION);
2699 let content_root = root.path().join("content");
2700 write_durable_post(&content_root);
2701 let durable_post_path = content_root.join("posts/durable-publication.md");
2702 fs::write(
2703 &durable_post_path,
2704 DURABLE_POST.replace(
2705 "slug = \"durable-publication\"\n",
2706 "slug = \"durable-publication\"\naliases = [\"original-durable-alias\"]\n",
2707 ),
2708 )
2709 .unwrap();
2710 let startup =
2711 StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
2712 let application = build_test_application(startup).await.unwrap();
2713 let admin = admin_client(&application);
2714 let response = admin_post_json(
2715 &admin,
2716 PUBLICATIONS_PATH,
2717 "bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb",
2718 &PublishNowRequest {
2719 post_id: uuid::Uuid::parse_str(DURABLE_POST_ID).unwrap(),
2720 preview_digest: admin_preview_digest(
2721 &admin,
2722 uuid::Uuid::parse_str(DURABLE_POST_ID).unwrap(),
2723 )
2724 .await,
2725 expected_revision: None,
2726 scheduled_for: None,
2727 },
2728 )
2729 .await;
2730 assert_eq!(response.status(), StatusCode::OK);
2731 let published: PublishNowResponse = admin_json(response).await;
2732 let public = reqwest::Client::builder()
2733 .no_proxy()
2734 .redirect(reqwest::redirect::Policy::none())
2735 .build()
2736 .unwrap();
2737 let initial_public_url = format!(
2738 "http://{}/posts/durable-publication",
2739 application.public_addr
2740 );
2741 let initial_alias_url = format!(
2742 "http://{}/posts/original-durable-alias",
2743 application.public_addr
2744 );
2745 let initial_alias = public.get(&initial_alias_url).send().await.unwrap();
2746 assert_eq!(initial_alias.status(), StatusCode::PERMANENT_REDIRECT);
2747 assert_eq!(
2748 initial_alias.headers()[reqwest::header::LOCATION],
2749 "https://startup.example.test/posts/durable-publication"
2750 );
2751 let initial = public.get(&initial_public_url).send().await.unwrap();
2752 let initial_etag = initial.headers()[reqwest::header::ETAG].clone();
2753 assert!(
2754 initial
2755 .text()
2756 .await
2757 .unwrap()
2758 .contains("Durable article body.")
2759 );
2760 stop_built_application(application).await;
2761
2762 let stopped_edit = DURABLE_POST
2763 .replace(
2764 "slug = \"durable-publication\"\n",
2765 "slug = \"durable-publication\"\naliases = [\"edited-durable-alias\"]\n",
2766 )
2767 .replace(
2768 "Durable article body.",
2769 "This edit was made while Maincopy was stopped.",
2770 );
2771 fs::write(&durable_post_path, stopped_edit).unwrap();
2772 let arguments = root.path().join("maincopy.toml");
2773 let startup =
2774 StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
2775 let restarted = build_test_application(startup).await.unwrap();
2776 let public_url = format!("http://{}/posts/durable-publication", restarted.public_addr);
2777 let pinned = public.get(&public_url).send().await.unwrap();
2778 assert_eq!(pinned.headers()[reqwest::header::ETAG], initial_etag);
2779 let pinned_body = pinned.text().await.unwrap();
2780 assert!(pinned_body.contains("Durable article body."));
2781 assert!(!pinned_body.contains("This edit was made while Maincopy was stopped."));
2782 let pinned_alias_url = format!(
2783 "http://{}/posts/original-durable-alias",
2784 restarted.public_addr
2785 );
2786 let pinned_alias = public.get(&pinned_alias_url).send().await.unwrap();
2787 assert_eq!(pinned_alias.status(), StatusCode::PERMANENT_REDIRECT);
2788 assert_eq!(
2789 pinned_alias.headers()[reqwest::header::LOCATION],
2790 "https://startup.example.test/posts/durable-publication"
2791 );
2792 let edited_alias_url = format!(
2793 "http://{}/posts/edited-durable-alias",
2794 restarted.public_addr
2795 );
2796 assert_eq!(
2797 public.get(&edited_alias_url).send().await.unwrap().status(),
2798 StatusCode::NOT_FOUND
2799 );
2800
2801 let admin = admin_client(&restarted);
2802 let posts: ListPostsResponse = admin_json(admin_get(&admin, POSTS_PATH).await).await;
2803 let edited_summary = posts
2804 .posts
2805 .into_iter()
2806 .find(|summary| summary.post_id.to_string() == DURABLE_POST_ID)
2807 .unwrap();
2808 assert_eq!(
2809 edited_summary.publication_state,
2810 PostPublicationState::UnpublishedChange
2811 );
2812 assert_ne!(edited_summary.revision, published.revision);
2813 let preview = admin_get(
2814 &admin,
2815 &format!("{POSTS_PATH}/{}/preview", edited_summary.post_id),
2816 )
2817 .await;
2818 assert_eq!(preview.status(), StatusCode::OK);
2819 let preview_body = admin_text(preview).await;
2820 assert!(preview_body.contains("This edit was made while Maincopy was stopped."));
2821 assert!(!preview_body.contains("Durable article body."));
2822
2823 let response = admin_post_json(
2824 &admin,
2825 PUBLICATIONS_PATH,
2826 "cccccccc-cccc-4ccc-8ccc-cccccccccccc",
2827 &PublishNowRequest {
2828 post_id: edited_summary.post_id,
2829 preview_digest: admin_preview_digest(&admin, edited_summary.post_id).await,
2830 expected_revision: Some(edited_summary.revision.clone()),
2831 scheduled_for: None,
2832 },
2833 )
2834 .await;
2835 assert_eq!(response.status(), StatusCode::OK);
2836 let approved: PublishNowResponse = admin_json(response).await;
2837 assert_eq!(approved.revision, edited_summary.revision);
2838
2839 let updated = public.get(&public_url).send().await.unwrap();
2840 assert_ne!(updated.headers()[reqwest::header::ETAG], initial_etag);
2841 let updated_body = updated.text().await.unwrap();
2842 assert!(updated_body.contains("This edit was made while Maincopy was stopped."));
2843 assert!(!updated_body.contains("Durable article body."));
2844 assert_eq!(
2845 public.get(&pinned_alias_url).send().await.unwrap().status(),
2846 StatusCode::NOT_FOUND
2847 );
2848 let updated_alias = public.get(&edited_alias_url).send().await.unwrap();
2849 assert_eq!(updated_alias.status(), StatusCode::PERMANENT_REDIRECT);
2850 assert_eq!(
2851 updated_alias.headers()[reqwest::header::LOCATION],
2852 "https://startup.example.test/posts/durable-publication"
2853 );
2854
2855 stop_built_application(restarted).await;
2856 }
2857
2858 #[tokio::test]
2859 #[cfg(target_os = "linux")]
2860 async fn restart_serves_the_durable_published_revision() {
2861 use maincopy_shared::{
2862 posts::{ListPostsResponse, POSTS_PATH, PostPublicationState},
2863 publication::{PUBLICATIONS_PATH, PublishNowRequest, PublishNowResponse},
2864 };
2865
2866 let (root, arguments, _) = startup_fixture("", VALID_PUBLICATION);
2867 let content_root = root.path().join("content");
2868 write_durable_post(&content_root);
2869 let startup =
2870 StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
2871 let application = build_test_application(startup).await.unwrap();
2872 let admin = admin_client(&application);
2873 let response = admin_post_json(
2874 &admin,
2875 PUBLICATIONS_PATH,
2876 "bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb",
2877 &PublishNowRequest {
2878 post_id: uuid::Uuid::parse_str(DURABLE_POST_ID).unwrap(),
2879 preview_digest: admin_preview_digest(
2880 &admin,
2881 uuid::Uuid::parse_str(DURABLE_POST_ID).unwrap(),
2882 )
2883 .await,
2884 expected_revision: None,
2885 scheduled_for: None,
2886 },
2887 )
2888 .await;
2889 assert_eq!(response.status(), StatusCode::OK);
2890 let published: PublishNowResponse = admin_json(response).await;
2891 let published_at = published.published_at.unwrap();
2892 let public = reqwest::Client::builder().no_proxy().build().unwrap();
2893 let initial_url = format!(
2894 "http://{}/posts/durable-publication",
2895 application.public_addr
2896 );
2897 let initial = public.get(&initial_url).send().await.unwrap();
2898 let initial_etag = initial.headers()[reqwest::header::ETAG].clone();
2899 let initial_html = initial.text().await.unwrap();
2900 assert!(initial_html.contains("Durable article body."));
2901 assert!(initial_html.contains(&published_at.to_string()));
2902
2903 stop_built_application(application).await;
2904
2905 let arguments = root.path().join("maincopy.toml");
2906 let startup =
2907 StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
2908 let application = build_test_application(startup).await.unwrap();
2909 let restarted_admin = admin_client(&application);
2910 let posts: ListPostsResponse =
2911 admin_json(admin_get(&restarted_admin, POSTS_PATH).await).await;
2912 let durable = posts
2913 .posts
2914 .into_iter()
2915 .find(|summary| summary.post_id == published.post_id)
2916 .unwrap();
2917 assert_eq!(durable.publication_state, PostPublicationState::Published);
2918 assert_eq!(durable.revision, published.revision);
2919
2920 let response = public
2921 .get(format!(
2922 "http://{}/posts/durable-publication",
2923 application.public_addr
2924 ))
2925 .send()
2926 .await
2927 .unwrap();
2928 assert_eq!(response.status(), reqwest::StatusCode::OK);
2929 assert_eq!(response.headers()[reqwest::header::ETAG], initial_etag);
2930 let html = response.text().await.unwrap();
2931 assert!(html.contains("Durable article body."));
2932 assert!(html.contains(&published_at.to_string()));
2933
2934 stop_built_application(application).await;
2935 }
2936
2937 #[tokio::test]
2938 #[cfg(target_os = "linux")]
2939 async fn offline_restore_requires_exact_scheduled_and_blocked_candidate_inputs() {
2940 use sqlx::{ConnectOptions as _, Connection as _};
2941
2942 let (root, arguments, _) = startup_fixture("", VALID_PUBLICATION);
2943 let content_root = root.path().join("content");
2944 write_durable_post(&content_root);
2945 let startup =
2946 StartupConfiguration::load_with_discovery(arguments.clone(), discover_content_tree)
2947 .unwrap();
2948 let database_path = startup._host.view().database.path.to_owned();
2949 let limits = startup._host.view().content_limits;
2950 let state_root = root.path().join("state");
2951 stop_built_application(build_test_application(startup).await.unwrap()).await;
2952
2953 let approved = discover_content_tree(&content_root, limits).unwrap();
2954 let approved_digest = approved.digest();
2955 let catalog =
2956 Arc::new(compile_content_catalog(&prepare_content(&approved).unwrap()).unwrap());
2957 let post_id = PostId::parse(DURABLE_POST_ID).unwrap();
2958 let revision = catalog.current_post(&post_id).unwrap().revision.clone();
2959 let preview = render_bound_post_preview(
2960 &catalog,
2961 embedded_manifest(),
2962 &post_id,
2963 None,
2964 "/api/admin/v1/preview-assets/retained-fixture",
2965 None,
2966 )
2967 .unwrap()
2968 .unwrap();
2969 let candidates = ContentCandidateStore::open(&state_root, limits).unwrap();
2970
2971 fs::write(
2974 content_root.join("publication.toml"),
2975 VALID_PUBLICATION.replace("Pinned startup source", "Another publication title"),
2976 )
2977 .unwrap();
2978 let same_revision = discover_content_tree(&content_root, limits).unwrap();
2979 assert_eq!(
2980 compile_content_catalog(&prepare_content(&same_revision).unwrap())
2981 .unwrap()
2982 .current_post(&post_id)
2983 .unwrap()
2984 .revision,
2985 revision,
2986 );
2987 let alternate_digest = candidates.retain(&same_revision).unwrap();
2988 assert_ne!(alternate_digest, approved_digest);
2989 fs::write(
2990 content_root.join("posts/durable-publication.md"),
2991 DURABLE_POST.replace("Durable article body.", "A later unapproved revision."),
2992 )
2993 .unwrap();
2994 let later = discover_content_tree(&content_root, limits).unwrap();
2995 candidates.retain(&later).unwrap();
2996
2997 let mut connection = sqlx::sqlite::SqliteConnectOptions::new()
2998 .filename(&database_path)
2999 .foreign_keys(true)
3000 .connect()
3001 .await
3002 .unwrap();
3003 sqlx::query(
3004 "INSERT INTO canonical_publications (\
3005 publication_id, creation_key, command_kind, stable_post_id, pinned_post_digest, \
3006 state, version, scheduled_at_ns, content_tree_digest, accepted_preview_digest\
3007 ) VALUES (?, ?, 'scheduled', ?, ?, 'scheduled', 1, ?, ?, ?)",
3008 )
3009 .bind(
3010 uuid::Uuid::parse_str(DURABLE_PUBLICATION_ID)
3011 .unwrap()
3012 .as_bytes()
3013 .as_slice(),
3014 )
3015 .bind(uuid::Uuid::new_v4().as_bytes().as_slice())
3016 .bind(post_id.as_uuid().as_bytes().as_slice())
3017 .bind(revision.as_bytes().as_slice())
3018 .bind(1_900_000_000_000_000_000_i64)
3019 .bind(approved_digest.as_bytes().as_slice())
3020 .bind(preview.digest.as_bytes().as_slice())
3021 .execute(&mut connection)
3022 .await
3023 .unwrap();
3024 connection.close().await.unwrap();
3025
3026 for (state, version, activation, reason) in [
3027 ("scheduled", 1_i64, None, None),
3028 (
3029 "blocked",
3030 3_i64,
3031 Some(1_900_000_000_000_000_000_i64),
3032 Some("preview_changed"),
3033 ),
3034 ] {
3035 let mut connection = sqlx::sqlite::SqliteConnectOptions::new()
3036 .filename(&database_path)
3037 .connect()
3038 .await
3039 .unwrap();
3040 sqlx::query(
3041 "UPDATE canonical_publications SET state = ?, version = ?, \
3042 activation_at_ns = ?, block_reason = ?",
3043 )
3044 .bind(state)
3045 .bind(version)
3046 .bind(activation)
3047 .bind(reason)
3048 .execute(&mut connection)
3049 .await
3050 .unwrap();
3051 connection.close().await.unwrap();
3052 verify_offline_content_without_mutation(&database_path, &state_root, limits)
3053 .await
3054 .unwrap();
3055
3056 fs::remove_file(
3057 state_root.join(format!("content-candidates/{approved_digest}.candidate")),
3058 )
3059 .unwrap();
3060 let error =
3061 verify_offline_content_without_mutation(&database_path, &state_root, limits)
3062 .await
3063 .unwrap_err();
3064 assert!(matches!(
3065 error,
3066 ProcessError::Application(ApplicationError::Startup {
3067 operation: "verify restored pinned candidates",
3068 ..
3069 })
3070 ));
3071 fs::remove_file(
3072 state_root.join(format!("content-candidates/{alternate_digest}.candidate")),
3073 )
3074 .unwrap();
3075 let error =
3076 verify_offline_content_without_mutation(&database_path, &state_root, limits)
3077 .await
3078 .unwrap_err();
3079 assert!(matches!(
3080 error,
3081 ProcessError::Application(ApplicationError::Startup {
3082 operation: "verify restored release revisions",
3083 ..
3084 })
3085 ));
3086 candidates.retain(&approved).unwrap();
3087 candidates.retain(&same_revision).unwrap();
3088 }
3089 }
3090
3091 async fn verify_offline_content_without_mutation(
3092 database_path: &Path,
3093 state_root: &Path,
3094 limits: ContentTreeLimits,
3095 ) -> Result<(), ProcessError> {
3096 let before = blake3::hash(&fs::read(database_path).unwrap());
3097 let inspected = database::restore::inspect(database_path).await.unwrap();
3098 let result = verify_restore_content(&inspected.store, state_root, limits).await;
3099 inspected.close().await;
3100 assert_eq!(blake3::hash(&fs::read(database_path).unwrap()), before);
3101 result
3102 }
3103
3104 #[tokio::test]
3105 #[cfg(target_os = "linux")]
3106 async fn startup_preflights_exact_activation_before_indexing_new_content_or_binding() {
3107 use sqlx::{ConnectOptions as _, Connection as _};
3108
3109 let (root, arguments, _) = startup_fixture("", VALID_PUBLICATION);
3110 let content_root = root.path().join("content");
3111 write_durable_post(&content_root);
3112 let startup =
3113 StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
3114 let database_path = startup._host.view().database.path.to_owned();
3115 let limits = startup._host.view().content_limits;
3116 let state_root = root.path().join("state");
3117 stop_built_application(build_test_application(startup).await.unwrap()).await;
3118
3119 let arguments = root.path().join("maincopy.toml");
3120 let startup =
3121 StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
3122 let content_digest = startup._content_tree.digest();
3123 let catalog = Arc::new(compile_content_catalog(&startup.content).unwrap());
3124 let post_id = PostId::parse(DURABLE_POST_ID).unwrap();
3125 let rendered = catalog.current_post(&post_id).unwrap();
3126 let revision = rendered.revision.clone();
3127 let accepted_preview_digest = render_bound_post_preview(
3128 &catalog,
3129 embedded_manifest(),
3130 &post_id,
3131 None,
3132 "/api/admin/v1/preview-assets/recovery-fixture",
3133 None,
3134 )
3135 .unwrap()
3136 .unwrap()
3137 .digest;
3138 let activation_at = OffsetDateTime::from_unix_timestamp(1_777_734_400).unwrap();
3139 let candidate_ledger = PublicLedgerProjection::empty()
3140 .with_published(PublishedPostRevision::new(
3141 post_id,
3142 revision.clone(),
3143 activation_at,
3144 ))
3145 .unwrap();
3146 let shell = render_site_shell(Arc::clone(&catalog), embedded_manifest(), &candidate_ledger)
3147 .unwrap();
3148 let candidate = shell.into_snapshot().unwrap();
3149 let candidate_digest = candidate.digest.clone();
3150 let activation_at_ns = i64::try_from(activation_at.unix_timestamp_nanos()).unwrap();
3151 let publication_id = uuid::Uuid::parse_str(DURABLE_PUBLICATION_ID)
3152 .unwrap()
3153 .into_bytes();
3154 let creation_key = uuid::Uuid::parse_str("bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb")
3155 .unwrap()
3156 .into_bytes();
3157 let stable_post_id = uuid::Uuid::parse_str(DURABLE_POST_ID).unwrap().into_bytes();
3158 let mut connection = sqlx::sqlite::SqliteConnectOptions::new()
3159 .filename(&database_path)
3160 .foreign_keys(true)
3161 .connect()
3162 .await
3163 .unwrap();
3164 sqlx::query(
3165 "INSERT INTO canonical_publications (\
3166 publication_id, creation_key, command_kind, stable_post_id, pinned_post_digest, state, version, \
3167 scheduled_at_ns, activation_at_ns, activation_site_digest, content_tree_digest, \
3168 accepted_preview_digest\
3169 ) VALUES (?, ?, 'immediate', ?, ?, 'activating', 2, ?, ?, ?, ?, ?)",
3170 )
3171 .bind(publication_id.as_slice())
3172 .bind(creation_key.as_slice())
3173 .bind(stable_post_id.as_slice())
3174 .bind(revision.as_bytes().as_slice())
3175 .bind(activation_at_ns)
3176 .bind(activation_at_ns)
3177 .bind(candidate_digest.as_bytes().as_slice())
3178 .bind(content_digest.as_bytes().as_slice())
3179 .bind(accepted_preview_digest.as_bytes().as_slice())
3180 .execute(&mut connection)
3181 .await
3182 .unwrap();
3183 let original_site: Vec<u8> =
3184 sqlx::query_scalar("SELECT current_site_digest FROM site_state WHERE singleton = 1")
3185 .fetch_one(&mut connection)
3186 .await
3187 .unwrap();
3188 sqlx::query("UPDATE canonical_publications SET activation_site_digest = ?")
3189 .bind([0xab_u8; 32].as_slice())
3190 .execute(&mut connection)
3191 .await
3192 .unwrap();
3193 connection.close().await.unwrap();
3194 drop(startup);
3195 let error = verify_offline_content_without_mutation(&database_path, &state_root, limits)
3196 .await
3197 .unwrap_err();
3198 assert!(matches!(
3199 error,
3200 ProcessError::Application(ApplicationError::Startup {
3201 stage: StartupStage::Content,
3202 operation: "verify restored activation",
3203 ..
3204 })
3205 ));
3206
3207 fs::write(
3208 content_root.join("posts/durable-publication.md"),
3209 DURABLE_POST.replace("Durable article body.", "Later private revision."),
3210 )
3211 .unwrap();
3212 let startup = StartupConfiguration::load_with_discovery(
3213 root.path().join("maincopy.toml"),
3214 discover_content_tree,
3215 )
3216 .unwrap();
3217 let error = match build_test_application(startup).await {
3218 Ok(_) => panic!("an unreproducible activation must reject startup"),
3219 Err(error) => error,
3220 };
3221 assert!(matches!(
3222 error,
3223 ProcessError::Application(ApplicationError::Startup {
3224 stage: StartupStage::Content,
3225 operation: "rebuild the activating publication candidate",
3226 ..
3227 })
3228 ));
3229
3230 let mut connection = sqlx::sqlite::SqliteConnectOptions::new()
3231 .filename(&database_path)
3232 .connect()
3233 .await
3234 .unwrap();
3235 let revisions: Vec<Vec<u8>> =
3236 sqlx::query_scalar("SELECT revision_digest FROM post_revisions")
3237 .fetch_all(&mut connection)
3238 .await
3239 .unwrap();
3240 assert_eq!(revisions, vec![revision.as_bytes().to_vec()]);
3241 let durable_site: Vec<u8> =
3242 sqlx::query_scalar("SELECT current_site_digest FROM site_state WHERE singleton = 1")
3243 .fetch_one(&mut connection)
3244 .await
3245 .unwrap();
3246 assert_eq!(durable_site, original_site);
3247 let state: String = sqlx::query_scalar("SELECT state FROM canonical_publications")
3248 .fetch_one(&mut connection)
3249 .await
3250 .unwrap();
3251 assert_eq!(state, "activating");
3252 sqlx::query("UPDATE canonical_publications SET activation_site_digest = ?")
3253 .bind(candidate_digest.as_bytes().as_slice())
3254 .execute(&mut connection)
3255 .await
3256 .unwrap();
3257 connection.close().await.unwrap();
3258 verify_offline_content_without_mutation(&database_path, &state_root, limits)
3259 .await
3260 .unwrap();
3261
3262 let startup = StartupConfiguration::load_with_discovery(
3263 root.path().join("maincopy.toml"),
3264 discover_content_tree,
3265 )
3266 .unwrap();
3267 let application = build_test_application(startup).await.unwrap();
3268 {
3269 let projection = application.publication_coordinator.read();
3270 assert_eq!(projection.site.digest, candidate_digest);
3271 assert_eq!(projection.ledger.len(), 1);
3272 let current = projection
3273 .catalog
3274 .current_post(&PostId::parse(DURABLE_POST_ID).unwrap())
3275 .unwrap();
3276 assert_ne!(current.revision, revision);
3277 assert!(
3278 current
3279 .article
3280 .identity_html
3281 .contains("Later private revision.")
3282 );
3283 }
3284 let response = reqwest::Client::builder()
3285 .no_proxy()
3286 .timeout(std::time::Duration::from_secs(2))
3287 .build()
3288 .unwrap()
3289 .get(format!(
3290 "http://{}/posts/durable-publication",
3291 application.public_addr
3292 ))
3293 .send()
3294 .await
3295 .unwrap();
3296 assert_eq!(response.status(), reqwest::StatusCode::OK);
3297 let html = response.text().await.unwrap();
3298 assert!(html.contains("Durable article body."));
3299 assert!(html.contains(&activation_at.to_string()));
3300 stop_built_application(application).await;
3301
3302 let mut connection = sqlx::sqlite::SqliteConnectOptions::new()
3303 .filename(&database_path)
3304 .connect()
3305 .await
3306 .unwrap();
3307 let (state, current_digest, published_at_ns): (String, Vec<u8>, i64) = sqlx::query_as(
3308 "SELECT state, current_published_digest, published_at_ns \
3309 FROM canonical_publications WHERE publication_id = ?",
3310 )
3311 .bind(publication_id.as_slice())
3312 .fetch_one(&mut connection)
3313 .await
3314 .unwrap();
3315 assert_eq!(state, "published");
3316 assert_eq!(current_digest, revision.as_bytes());
3317 assert_eq!(published_at_ns, activation_at_ns);
3318 connection.close().await.unwrap();
3319 }
3320
3321 #[tokio::test]
3322 #[cfg(target_os = "linux")]
3323 async fn mail_credential_failure_releases_database_process_and_listener_ownership() {
3324 use std::{io::Write as _, os::unix::fs::OpenOptionsExt as _};
3325
3326 let (root, config_path, _) = startup_fixture("", VALID_PUBLICATION);
3327 let reservations = [
3328 reserve_loopback_port(),
3329 reserve_loopback_port(),
3330 reserve_loopback_port(),
3331 ];
3332 let addresses = reservations
3333 .each_ref()
3334 .map(|socket| socket.local_addr().unwrap());
3335 let host_source = startup_host_source_with_listeners(
3336 "",
3337 &addresses[0].to_string(),
3338 &addresses[1].to_string(),
3339 &addresses[2].to_string(),
3340 );
3341 fs::write(
3342 &config_path,
3343 format!(
3344 "{host_source}\n\
3345 [mail]\n\
3346 mode = \"ses\"\n\
3347 sender = \"newsletter@example.com\"\n\
3348 region = \"us-east-1\"\n\
3349 configuration_set = \"newsletter\"\n\
3350 credential_file = \"private-mail-credential.json\"\n\
3351 control_signing_key_file = \"unused-control.key\"\n"
3352 ),
3353 )
3354 .unwrap();
3355 let credential_path = root.path().join("private-mail-credential.json");
3356 let mut credential = fs::OpenOptions::new()
3357 .write(true)
3358 .create_new(true)
3359 .mode(0o600)
3360 .open(&credential_path)
3361 .unwrap();
3362 credential
3365 .write_all(br#"{"access_key_id":"AKIDEXAMPLE","secret_access_key":"private-fixture-secret-material","unexpected":"invalid"}"#)
3366 .unwrap();
3367 drop(credential);
3368 let startup =
3369 StartupConfiguration::load_with_discovery(config_path.clone(), discover_content_tree)
3370 .unwrap();
3371 bootstrap_test_identity(&startup).await;
3372 let error = match Application::build(startup).await {
3373 Ok(_) => panic!("malformed SES credentials must fail application construction"),
3374 Err(error) => error,
3375 };
3376 assert!(matches!(
3377 &error,
3378 ProcessError::Application(ApplicationError::Startup {
3379 stage: StartupStage::Configuration,
3380 operation: "prepare mail runtime",
3381 source,
3382 }) if matches!(source.downcast_ref::<MailStartupError>(),
3383 Some(MailStartupError::Review(MailReviewAccessError::CredentialInvalid)))
3384 ));
3385 let diagnostic = format!("{error}\n{error:?}");
3386 for private in [
3387 credential_path.to_str().unwrap(),
3388 "private-mail-credential.json",
3389 "private-fixture-secret-material",
3390 ] {
3391 assert!(!diagnostic.contains(private));
3392 }
3393
3394 let configuration = HostConfigurationLoader::from_process_working_directory()
3397 .unwrap()
3398 .load(&config_path)
3399 .unwrap();
3400 let host = configuration.view();
3401 let process_lock = ProcessLock::acquire(host.runtime_root).unwrap();
3402 let database = database::bootstrap(host.database).await.unwrap();
3403 for address in addresses {
3404 let listener = tokio::net::TcpListener::bind(address).await.unwrap();
3405 drop(listener);
3406 }
3407 database.close().await.unwrap();
3408 drop(process_lock);
3409 drop(reservations);
3410 }
3411
3412 #[tokio::test]
3413 #[cfg(target_os = "linux")]
3414 async fn listener_failure_releases_the_public_port_and_database_ownership() {
3415 let (root, _, _) = startup_fixture("", VALID_PUBLICATION);
3416 let config_path = root.path().join("maincopy.toml");
3417 let reservation = reserve_loopback_port();
3418 let public_addr = reservation.local_addr().unwrap();
3419 let occupied = tokio::net::TcpListener::bind(public_addr).await.unwrap();
3420 let public_bind = public_addr.to_string();
3421 fs::write(&config_path, startup_host_source("", &public_bind)).unwrap();
3422 let arguments = config_path.clone();
3423 let startup =
3424 StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
3425 let error = match build_test_application(startup).await {
3426 Ok(_) => panic!("an occupied public address must fail listener binding"),
3427 Err(error) => error,
3428 };
3429
3430 assert!(matches!(
3431 error,
3432 ProcessError::Application(ApplicationError::Startup {
3433 stage: StartupStage::Listeners,
3434 ..
3435 })
3436 ));
3437 drop(occupied);
3438 let arguments = config_path;
3439 let startup =
3440 StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
3441 let application = build_test_application(startup).await.unwrap();
3442 assert_eq!(application.public_addr, public_addr);
3443
3444 stop_built_application(application).await;
3445 drop(reservation);
3446 }
3447
3448 #[tokio::test]
3449 #[cfg(target_os = "linux")]
3450 async fn admin_listener_failure_releases_public_listener_and_database_ownership() {
3451 let (root, _, _) = startup_fixture("", VALID_PUBLICATION);
3452 let config_path = root.path().join("maincopy.toml");
3453 let reservation = reserve_loopback_port();
3454 let admin_addr = reservation.local_addr().unwrap();
3455 let occupied = tokio::net::TcpListener::bind(admin_addr).await.unwrap();
3456 fs::write(
3457 &config_path,
3458 startup_host_source_with_listeners(
3459 "",
3460 "127.0.0.1:0",
3461 &admin_addr.to_string(),
3462 "127.0.0.1:0",
3463 ),
3464 )
3465 .unwrap();
3466 let arguments = config_path.clone();
3467 let startup =
3468 StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
3469 let error = match build_test_application(startup).await {
3470 Ok(_) => panic!("an occupied admin address must fail listener binding"),
3471 Err(error) => error,
3472 };
3473
3474 assert!(matches!(
3475 error,
3476 ProcessError::Application(ApplicationError::Startup {
3477 stage: StartupStage::Listeners,
3478 ..
3479 })
3480 ));
3481 drop(occupied);
3482 let arguments = config_path;
3483 let startup =
3484 StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
3485 let application = build_test_application(startup).await.unwrap();
3486 assert_eq!(application.admin_addr, admin_addr);
3487
3488 stop_built_application(application).await;
3489 drop(reservation);
3490 }
3491
3492 #[tokio::test]
3493 #[cfg(target_os = "linux")]
3494 async fn database_failure_prevents_application_build() {
3495 use std::os::unix::fs::PermissionsExt as _;
3496
3497 use sqlx::{ConnectOptions as _, Connection as _};
3498
3499 let (_root, arguments, _) = startup_fixture("", VALID_PUBLICATION);
3500 let startup =
3501 StartupConfiguration::load_with_discovery(arguments, discover_content_tree).unwrap();
3502 let database_path = startup._host.view().database.path.to_owned();
3503 let database_parent = database_path.parent().unwrap();
3504 fs::create_dir_all(database_parent).unwrap();
3505 fs::set_permissions(database_parent, fs::Permissions::from_mode(0o700)).unwrap();
3506
3507 let mut foreign = sqlx::sqlite::SqliteConnectOptions::new()
3508 .filename(&database_path)
3509 .create_if_missing(true)
3510 .connect()
3511 .await
3512 .unwrap();
3513 sqlx::query("PRAGMA application_id = 7")
3514 .execute(&mut foreign)
3515 .await
3516 .unwrap();
3517 foreign.close().await.unwrap();
3518 fs::set_permissions(&database_path, fs::Permissions::from_mode(0o600)).unwrap();
3519
3520 let error = match Application::build(startup).await {
3521 Ok(_) => panic!("a foreign database must fail before listener binding"),
3522 Err(error) => error,
3523 };
3524
3525 assert!(matches!(
3526 error,
3527 ProcessError::Application(ApplicationError::Startup {
3528 stage: StartupStage::Database,
3529 ..
3530 })
3531 ));
3532 }
3533
3534 async fn begin_public_request(address: std::net::SocketAddr, path: &str) -> TcpStream {
3535 let mut stream = TcpStream::connect(address).await.unwrap();
3536 let head = format!(
3537 "GET {path} HTTP/1.1\r\nHost: localhost\r\nContent-Length: 1\r\nExpect: 100-continue\r\nConnection: close\r\n\r\n"
3538 );
3539 stream.write_all(head.as_bytes()).await.unwrap();
3540 let mut interim = Vec::new();
3541 tokio::time::timeout(std::time::Duration::from_secs(2), async {
3542 while !interim.ends_with(b"\r\n\r\n") {
3543 assert!(interim.len() < 1024);
3544 interim.push(stream.read_u8().await.unwrap());
3545 }
3546 })
3547 .await
3548 .unwrap();
3549 assert_eq!(
3550 std::str::from_utf8(&interim).unwrap(),
3551 "HTTP/1.1 100 Continue\r\n\r\n"
3552 );
3553 stream
3554 }
3555
3556 async fn finish_public_request(mut stream: TcpStream) -> String {
3557 stream.write_all(b"x").await.unwrap();
3558 let mut response = String::new();
3559 tokio::time::timeout(
3560 std::time::Duration::from_secs(2),
3561 stream.take(64 * 1024).read_to_string(&mut response),
3562 )
3563 .await
3564 .unwrap()
3565 .unwrap();
3566 response
3567 }
3568
3569 #[tokio::test]
3570 async fn shutdown_finishes_an_accepted_public_request_before_closing_the_real_writer() {
3571 let (_root, config, _) = startup_fixture("", VALID_PUBLICATION);
3572 let reservation = reserve_loopback_port();
3573 let address = reservation.local_addr().unwrap();
3574 fs::write(&config, startup_host_source("", &address.to_string())).unwrap();
3575 let startup =
3576 StartupConfiguration::load_with_discovery(config, discover_content_tree).unwrap();
3577 let mut application = build_test_application(startup).await.unwrap();
3578 assert_eq!(application.public_addr, address);
3579 let readiness = application.runtime.readiness.clone();
3580 let cancellation = application.runtime.cancellation.clone();
3581 let database_shutdown = application.runtime.database_shutdown.clone();
3582 let (stop, stopped) = oneshot::channel();
3583 application.runtime.shutdown = Box::pin(async {
3584 stopped.await.unwrap();
3585 Ok(())
3586 });
3587 let running = tokio::spawn(application.run_until_stop());
3588 let request = begin_public_request(address, "/health/live").await;
3589 stop.send(()).unwrap();
3590 tokio::time::timeout(std::time::Duration::from_secs(2), cancellation.cancelled())
3591 .await
3592 .unwrap();
3593 assert!(!readiness.is_ready());
3594 assert!(!database_shutdown.is_cancelled());
3595 assert!(!running.is_finished());
3596 let response = finish_public_request(request).await;
3597 assert!(response.starts_with("HTTP/1.1 200 OK"));
3598 assert!(response.contains(r#"{"status":"live"}"#));
3599 tokio::time::timeout(std::time::Duration::from_secs(2), running)
3600 .await
3601 .unwrap()
3602 .unwrap()
3603 .unwrap();
3604 assert!(database_shutdown.is_cancelled());
3605 let rebound = tokio::net::TcpListener::bind(address).await.unwrap();
3606 assert_eq!(rebound.local_addr().unwrap(), address);
3607 drop(rebound);
3608 drop(reservation);
3609 }
3610
3611 #[tokio::test]
3612 async fn supervised_task_failure_makes_pending_readiness_fail_while_liveness_stays_live() {
3613 let (_root, config, _) = startup_fixture("", VALID_PUBLICATION);
3614 let startup =
3615 StartupConfiguration::load_with_discovery(config, discover_content_tree).unwrap();
3616 let mut application = build_test_application(startup).await.unwrap();
3617 let address = application.public_addr;
3618 let cancellation = application.runtime.cancellation.clone();
3619 let (fail, failed) = oneshot::channel();
3620 application.runtime.critical_tasks.spawn(CriticalTask::new(
3621 CriticalTaskName::Worker,
3622 async {
3623 failed.await.unwrap();
3624 Ok::<(), CriticalTaskFailure>(())
3625 },
3626 ));
3627 let running = tokio::spawn(application.run_until_stop());
3628 let ready = begin_public_request(address, "/health/ready").await;
3629 let live = begin_public_request(address, "/health/live").await;
3630 fail.send(()).unwrap();
3631 tokio::time::timeout(std::time::Duration::from_secs(2), cancellation.cancelled())
3632 .await
3633 .unwrap();
3634 let ready = finish_public_request(ready).await;
3635 let live = finish_public_request(live).await;
3636 assert!(ready.starts_with("HTTP/1.1 503 Service Unavailable"));
3637 assert!(ready.contains(r#"{"status":"not_ready"}"#));
3638 assert!(live.starts_with("HTTP/1.1 200 OK"));
3639 assert!(live.contains(r#"{"status":"live"}"#));
3640 let failure = tokio::time::timeout(std::time::Duration::from_secs(2), running)
3641 .await
3642 .unwrap()
3643 .unwrap()
3644 .unwrap_err();
3645 assert!(matches!(
3646 failure,
3647 ApplicationError::CriticalTaskExited {
3648 task: CriticalTaskName::Worker
3649 }
3650 ));
3651 }
3652
3653 #[tokio::test]
3654 async fn shutdown_signal_marks_unready_cancels_and_drains_every_task() {
3655 let readiness = Readiness::default();
3656 let cancellation = CancellationToken::new();
3657 let drained = Arc::new(AtomicUsize::new(0));
3658 let observed_order = Arc::new(AtomicUsize::new(0));
3659 let (shutdown_tx, shutdown_rx) = oneshot::channel();
3660 let tasks = (0..2)
3661 .map(|_| {
3662 let readiness = readiness.clone();
3663 let cancellation = cancellation.clone();
3664 let drained = Arc::clone(&drained);
3665 let observed_order = Arc::clone(&observed_order);
3666 CriticalTask::new(CriticalTaskName::Worker, async move {
3667 cancellation.cancelled().await;
3668 if !readiness.is_ready() {
3669 observed_order.fetch_add(1, Ordering::SeqCst);
3670 }
3671 drained.fetch_add(1, Ordering::SeqCst);
3672 Ok::<(), CriticalTaskFailure>(())
3673 })
3674 })
3675 .collect();
3676 let application = ApplicationRuntime::with_parts(
3677 readiness.clone(),
3678 cancellation.clone(),
3679 Box::pin(async move {
3680 let _ = shutdown_rx.await;
3681 Ok(())
3682 }),
3683 tasks,
3684 );
3685
3686 let running = tokio::spawn(application.run_until_stop());
3687 while !readiness.is_ready() {
3688 tokio::task::yield_now().await;
3689 }
3690 assert!(readiness.is_ready());
3691
3692 assert!(shutdown_tx.send(()).is_ok());
3693 assert!(running.await.unwrap().is_ok());
3694 assert!(!readiness.is_ready());
3695 assert!(cancellation.is_cancelled());
3696 assert_eq!(drained.load(Ordering::SeqCst), 2);
3697 assert_eq!(observed_order.load(Ordering::SeqCst), 2);
3698 }
3699
3700 #[tokio::test]
3701 async fn shutdown_waits_for_the_database_writer_to_close() {
3702 let readiness = Readiness::default();
3703 let cancellation = CancellationToken::new();
3704 let database_shutdown = CancellationToken::new();
3705 let writer_shutdown = database_shutdown.clone();
3706 let (shutdown_tx, shutdown_rx) = oneshot::channel();
3707 let (writer_started_tx, writer_started_rx) = oneshot::channel();
3708 let (writer_release_tx, writer_release_rx) = oneshot::channel();
3709 let writer = CriticalTask::new(CriticalTaskName::DatabaseWriter, async move {
3710 writer_shutdown.cancelled().await;
3711 let _ = writer_started_tx.send(());
3712 let _ = writer_release_rx.await;
3713 Ok::<(), CriticalTaskFailure>(())
3714 });
3715 let application = ApplicationRuntime::with_database_writer(
3716 readiness.clone(),
3717 cancellation,
3718 database_shutdown,
3719 Box::pin(async move {
3720 let _ = shutdown_rx.await;
3721 Ok(())
3722 }),
3723 Vec::new(),
3724 spawn_critical_task(writer),
3725 );
3726
3727 let running = tokio::spawn(application.run_until_stop());
3728 while !readiness.is_ready() {
3729 tokio::task::yield_now().await;
3730 }
3731
3732 assert!(shutdown_tx.send(()).is_ok());
3733 assert!(writer_started_rx.await.is_ok());
3734 tokio::task::yield_now().await;
3735 assert!(!running.is_finished());
3736
3737 assert!(writer_release_tx.send(()).is_ok());
3738 assert!(running.await.unwrap().is_ok());
3739 assert!(!readiness.is_ready());
3740 }
3741
3742 #[tokio::test]
3743 async fn shutdown_drains_producers_before_stopping_the_database_writer() {
3744 let readiness = Readiness::default();
3745 let cancellation = CancellationToken::new();
3746 let producer_shutdown = cancellation.clone();
3747 let database_shutdown = CancellationToken::new();
3748 let writer_shutdown = database_shutdown.clone();
3749 let (shutdown_tx, shutdown_rx) = oneshot::channel();
3750 let (producer_cancelled_tx, producer_cancelled_rx) = oneshot::channel();
3751 let (producer_release_tx, producer_release_rx) = oneshot::channel();
3752 let (writer_cancelled_tx, mut writer_cancelled_rx) = oneshot::channel();
3753 let producer = CriticalTask::new(CriticalTaskName::Worker, async move {
3754 producer_shutdown.cancelled().await;
3755 let _ = producer_cancelled_tx.send(());
3756 let _ = producer_release_rx.await;
3757 Ok::<(), CriticalTaskFailure>(())
3758 });
3759 let writer = CriticalTask::new(CriticalTaskName::DatabaseWriter, async move {
3760 writer_shutdown.cancelled().await;
3761 let _ = writer_cancelled_tx.send(());
3762 Ok::<(), CriticalTaskFailure>(())
3763 });
3764 let application = ApplicationRuntime::with_database_writer(
3765 readiness.clone(),
3766 cancellation,
3767 database_shutdown,
3768 Box::pin(async move {
3769 let _ = shutdown_rx.await;
3770 Ok(())
3771 }),
3772 vec![producer],
3773 spawn_critical_task(writer),
3774 );
3775
3776 let running = tokio::spawn(application.run_until_stop());
3777 while !readiness.is_ready() {
3778 tokio::task::yield_now().await;
3779 }
3780 assert!(shutdown_tx.send(()).is_ok());
3781 assert!(producer_cancelled_rx.await.is_ok());
3782 assert!(matches!(
3783 writer_cancelled_rx.try_recv(),
3784 Err(oneshot::error::TryRecvError::Empty)
3785 ));
3786
3787 assert!(producer_release_tx.send(()).is_ok());
3788 assert!(writer_cancelled_rx.await.is_ok());
3789 assert!(running.await.unwrap().is_ok());
3790 }
3791
3792 #[tokio::test]
3793 async fn unexpected_success_cancels_and_drains_before_returning_failure() {
3794 let (result, readiness, cancellation, companion_drained) =
3795 run_with_unexpected_task(async { Ok(()) }).await;
3796
3797 assert!(matches!(
3798 result,
3799 Err(ApplicationError::CriticalTaskExited {
3800 task: CriticalTaskName::Scheduler
3801 })
3802 ));
3803 assert_shutdown_state(readiness, cancellation, companion_drained);
3804 }
3805
3806 #[tokio::test]
3807 async fn task_error_cancels_and_drains_before_returning_failure() {
3808 let (result, readiness, cancellation, companion_drained) =
3809 run_with_unexpected_task(async {
3810 Err(Box::new(std::io::Error::other("task failed")) as CriticalTaskFailure)
3811 })
3812 .await;
3813
3814 assert!(matches!(
3815 result,
3816 Err(ApplicationError::CriticalTaskFailed {
3817 task: CriticalTaskName::Scheduler,
3818 ..
3819 })
3820 ));
3821 assert_shutdown_state(readiness, cancellation, companion_drained);
3822 }
3823
3824 #[tokio::test]
3825 async fn task_panic_cancels_and_drains_before_returning_failure() {
3826 async fn panicking_task() -> Result<(), CriticalTaskFailure> {
3827 panic!("task panicked")
3828 }
3829
3830 let (result, readiness, cancellation, companion_drained) =
3831 run_with_unexpected_task(panicking_task()).await;
3832
3833 assert!(matches!(
3834 result,
3835 Err(ApplicationError::CriticalTaskPanicked {
3836 task: CriticalTaskName::Scheduler,
3837 ..
3838 })
3839 ));
3840 assert_shutdown_state(readiness, cancellation, companion_drained);
3841 }
3842
3843 #[tokio::test]
3844 async fn metrics_collector_storage_failure_marks_unready_and_drains_the_writer() {
3845 let (root, config_path, _publication_path) = startup_fixture("", VALID_PUBLICATION);
3846 let startup =
3847 StartupConfiguration::load_with_discovery(config_path, discover_content_tree).unwrap();
3848 let host = startup._host.view();
3849 let bootstrapped = database::bootstrap(host.database).await.unwrap();
3850 let (store, writer) = bootstrapped.into_store(4);
3851 let metrics = Metrics::new(&store.health.metrics).unwrap();
3852 let collector = MetricsCollector::new(
3853 metrics,
3854 store.health.clone(),
3855 BackupHealth::new(None),
3856 tokio::runtime::Handle::current(),
3857 );
3858 fs::rename(root.path().join("state"), root.path().join("moved-state")).unwrap();
3860 fs::write(root.path().join("state"), []).unwrap();
3861 let readiness = Readiness::default();
3862 let cancellation = CancellationToken::new();
3863 let writer_shutdown = CancellationToken::new();
3864 let application = ApplicationRuntime::with_database_writer(
3865 readiness.clone(),
3866 cancellation.clone(),
3867 writer_shutdown.clone(),
3868 Box::pin(std::future::pending()),
3869 vec![CriticalTask::new(
3870 CriticalTaskName::MetricsCollector,
3871 collector.run(cancellation.clone()),
3872 )],
3873 spawn_critical_task(CriticalTask::new(
3874 CriticalTaskName::DatabaseWriter,
3875 writer.run(writer_shutdown.clone()),
3876 )),
3877 );
3878 let result = tokio::time::timeout(
3879 std::time::Duration::from_secs(5),
3880 application.run_until_stop(),
3881 )
3882 .await
3883 .unwrap();
3884 assert!(matches!(
3885 result,
3886 Err(ApplicationError::CriticalTaskFailed {
3887 task: CriticalTaskName::MetricsCollector,
3888 ..
3889 })
3890 ));
3891 assert!(!readiness.is_ready());
3892 assert!(cancellation.is_cancelled());
3893 assert!(writer_shutdown.is_cancelled());
3894 assert_eq!(store.health.metrics.writer_up.get(), 0);
3895 }
3896
3897 #[tokio::test]
3898 async fn database_writer_failure_marks_unready_and_drains_producers() {
3899 let readiness = Readiness::default();
3900 let cancellation = CancellationToken::new();
3901 let database_shutdown = CancellationToken::new();
3902 let companion_drained = Arc::new(AtomicBool::new(false));
3903 let companion = {
3904 let cancellation = cancellation.clone();
3905 let companion_drained = Arc::clone(&companion_drained);
3906 CriticalTask::new(CriticalTaskName::Worker, async move {
3907 cancellation.cancelled().await;
3908 companion_drained.store(true, Ordering::SeqCst);
3909 Ok::<(), CriticalTaskFailure>(())
3910 })
3911 };
3912 let writer = CriticalTask::new(CriticalTaskName::DatabaseWriter, async {
3913 Err(std::io::Error::other("writer stopped"))
3914 });
3915 let application = ApplicationRuntime::with_database_writer(
3916 readiness.clone(),
3917 cancellation.clone(),
3918 database_shutdown.clone(),
3919 Box::pin(std::future::pending()),
3920 vec![companion],
3921 spawn_critical_task(writer),
3922 );
3923
3924 let result = application.run_until_stop().await;
3925
3926 assert!(matches!(
3927 result,
3928 Err(ApplicationError::CriticalTaskFailed {
3929 task: CriticalTaskName::DatabaseWriter,
3930 source,
3931 }) if source.downcast_ref::<std::io::Error>()
3932 .is_some_and(|error| error.to_string() == "writer stopped")
3933 ));
3934 assert!(!readiness.is_ready());
3935 assert!(cancellation.is_cancelled());
3936 assert!(database_shutdown.is_cancelled());
3937 assert!(companion_drained.load(Ordering::SeqCst));
3938 }
3939
3940 async fn run_with_unexpected_task<Future>(
3941 trigger: Future,
3942 ) -> (
3943 Result<(), ApplicationError>,
3944 Readiness,
3945 CancellationToken,
3946 Arc<AtomicBool>,
3947 )
3948 where
3949 Future: std::future::Future<Output = CriticalTaskResult> + Send + 'static,
3950 {
3951 let readiness = Readiness::default();
3952 let cancellation = CancellationToken::new();
3953 let companion_drained = Arc::new(AtomicBool::new(false));
3954 let companion = {
3955 let cancellation = cancellation.clone();
3956 let companion_drained = Arc::clone(&companion_drained);
3957 async move {
3958 cancellation.cancelled().await;
3959 companion_drained.store(true, Ordering::SeqCst);
3960 Ok::<(), CriticalTaskFailure>(())
3961 }
3962 };
3963 let application = ApplicationRuntime::with_parts(
3964 readiness.clone(),
3965 cancellation.clone(),
3966 Box::pin(std::future::pending()),
3967 vec![
3968 CriticalTask::new(CriticalTaskName::Scheduler, trigger),
3969 CriticalTask::new(CriticalTaskName::Worker, companion),
3970 ],
3971 );
3972
3973 let result = application.run_until_stop().await;
3974 (result, readiness, cancellation, companion_drained)
3975 }
3976
3977 fn assert_shutdown_state(
3978 readiness: Readiness,
3979 cancellation: CancellationToken,
3980 companion_drained: Arc<AtomicBool>,
3981 ) {
3982 assert!(!readiness.is_ready());
3983 assert!(cancellation.is_cancelled());
3984 assert!(companion_drained.load(Ordering::SeqCst));
3985 }
3986}