use gwk_domain::blob::BlobAddress;
use gwk_domain::ids::EvidenceId;
use gwk_domain::port::BlobStore as _;
use gwk_domain::protocol::KernelErrorCode;
use gwk_kernel::config::{AdminConfig, BlobConfig, KernelConfig};
use gwk_kernel::project::Refusal;
use gwk_kernel::wire::listen::{Listener, notify_ready};
use gwk_kernel::wire::serve::{self, Daemon};
use gwk_kernel::{InitOutcome, PgBlobStore, PgEventStore, TargetState, WriterLock, admin, recover};
use secrecy::{ExposeSecret as _, SecretString};
use serde_json::{Value, json};
use crate::exit::Failure;
use crate::{PUBLIC_REVISION, emit};
const POOL_CONNECTIONS: u32 = (gwk_kernel::MAX_INFLIGHT_APPENDS as u32) * 2;
#[derive(Debug, PartialEq, Eq)]
pub enum Retention {
Pin {
address: BlobAddress,
evidence: String,
},
Unpin {
address: BlobAddress,
evidence: String,
},
Sweep,
Shred {
address: BlobAddress,
},
}
pub async fn daemon(pretty: bool) -> Result<(), Failure> {
let revision = revision()?;
let config = KernelConfig::from_env().map_err(configuration)?;
let blob_config = BlobConfig::from_env().map_err(configuration)?;
let lock = WriterLock::acquire(config.database_url())
.await
.map_err(|e| Failure::new(KernelErrorCode::Fenced, e.to_string()))?;
let pool = gwk_kernel::connect_pool(config.database_url(), POOL_CONNECTIONS)
.await
.map_err(configuration)?;
let privileges = admin::runtime_privileges(&pool)
.await
.map_err(configuration)?;
let violations = privileges.violations();
if !violations.is_empty() {
return Err(Failure::new(
KernelErrorCode::Privilege,
format!(
"this credential holds privileges the kernel refuses to run with: {}",
violations.join(", ")
),
));
}
let blobs = PgBlobStore::open(pool.clone(), blob_config)
.await
.map_err(blob_failure)?;
let store = PgEventStore::open(pool)
.await
.map_err(configuration)?
.with_blobs(blobs);
let recovered = store.recover().await.map_err(refusal)?;
if !recovered.ready() {
return Err(Failure::new(
KernelErrorCode::Storage,
format!(
"the projections do not agree with the log ({:?}); refusing to serve",
recovered.verdict
),
));
}
let daemon =
std::sync::Arc::new(Daemon::new(store, revision.to_owned()).map_err(configuration)?);
let listener = Listener::bind(config.socket_path())
.await
.map_err(configuration)?;
emit(
&json!({
"type": "daemon_started",
"socket_path": config.socket_path().to_string_lossy(),
"public_revision": revision,
"watermark": recovered.watermark.map(|seq| seq.value().to_string()),
"verdict": verdict(&recovered),
"rejected_checkpoints": recovered.rejected.len(),
"uncertain_attempts": recovered.uncertain,
"notified_systemd": notify_ready(),
}),
pretty,
);
let stopped = serve::run(listener, daemon, async move {
tokio::select! {
_ = tokio::signal::ctrl_c() => {}
() = terminated() => {}
() = lock.cancelled() => {}
}
})
.await
.map_err(configuration)?;
emit(
&json!({
"type": "daemon_stopped",
"checkpoint": stopped.checkpoint.map(|seq| seq.value().to_string()),
"checkpoint_error": stopped.checkpoint_error,
}),
pretty,
);
Ok(())
}
async fn terminated() {
use tokio::signal::unix::{SignalKind, signal};
match signal(SignalKind::terminate()) {
Ok(mut term) => {
term.recv().await;
}
Err(_) => std::future::pending().await,
}
}
pub async fn init(pretty: bool) -> Result<(), Failure> {
let revision = revision()?;
let config = AdminConfig::from_env().map_err(configuration)?;
let _lock = WriterLock::acquire(config.admin_database_url())
.await
.map_err(|e| Failure::new(KernelErrorCode::Fenced, e.to_string()))?;
let pool = gwk_kernel::connect_pool(config.admin_database_url(), 4)
.await
.map_err(configuration)?;
let outcome = admin::init(&pool, &config).await.map_err(configuration)?;
let store = PgEventStore::open(pool).await.map_err(configuration)?;
store.ensure_genesis(&revision).await.map_err(refusal)?;
emit(
&json!({
"type": "admin_initialized",
"outcome": match outcome {
InitOutcome::Initialized => "initialized",
InitOutcome::AlreadyInitialized => "already_initialized",
},
"runtime_role": config.runtime_role(),
"public_revision": revision,
"contract_sha256": gwk_kernel::CONTRACT_SQL_SHA256,
}),
pretty,
);
Ok(())
}
pub async fn verify(pretty: bool) -> Result<(), Failure> {
let config = AdminConfig::from_env().map_err(configuration)?;
let pool = gwk_kernel::connect_pool(config.admin_database_url(), 2)
.await
.map_err(configuration)?;
let state = admin::inspect(&pool).await.map_err(configuration)?;
let attributes = admin::role_attributes(&pool, config.runtime_role())
.await
.map_err(configuration)?;
let violations: Vec<&'static str> = attributes
.map(|attributes| attributes.violations())
.unwrap_or_default();
let (target, detail) = match &state {
TargetState::Empty => ("empty", Value::Null),
TargetState::Initialized { contract_sha256 } => {
("initialized", json!({"contract_sha256": contract_sha256}))
}
TargetState::Foreign { objects } => ("foreign", json!({"objects": objects})),
};
emit(
&json!({
"type": "admin_verified",
"target": target,
"detail": detail,
"runtime_role": config.runtime_role(),
"runtime_role_exists": attributes.is_some(),
"violations": violations,
"expected_contract_sha256": gwk_kernel::CONTRACT_SQL_SHA256,
}),
pretty,
);
if !violations.is_empty() {
return Err(Failure::new(
KernelErrorCode::Privilege,
format!("the runtime role holds {}", violations.join(", ")),
));
}
if let TargetState::Initialized { contract_sha256 } = &state
&& contract_sha256 != gwk_kernel::CONTRACT_SQL_SHA256
{
return Err(Failure::new(
KernelErrorCode::Schema,
format!(
"the database carries contract {contract_sha256}, and this build is {}",
gwk_kernel::CONTRACT_SQL_SHA256
),
));
}
Ok(())
}
pub async fn rebuild_projections(scratch: &str, pretty: bool) -> Result<(), Failure> {
let config = AdminConfig::from_env().map_err(configuration)?;
let live = gwk_kernel::connect_pool(config.admin_database_url(), 4)
.await
.map_err(configuration)?;
let scratch_url = beside(config.admin_database_url(), scratch)?;
let scratch_pool = gwk_kernel::connect_pool(&scratch_url, 4)
.await
.map_err(configuration)?;
let scratch_config = AdminConfig::from_lookup({
let url = scratch_url.expose_secret().to_owned();
let role = config.runtime_role().to_owned();
move |key| match key {
gwk_kernel::config::ADMIN_DATABASE_URL_ENV => Some(url.clone()),
gwk_kernel::config::RUNTIME_ROLE_ENV => Some(role.clone()),
_ => None,
}
})
.map_err(configuration)?;
admin::init(&scratch_pool, &scratch_config)
.await
.map_err(configuration)?;
let report = PgEventStore::open_reader(live)
.rebuild_into(&scratch_pool)
.await
.map_err(refusal)?;
emit(
&json!({
"type": "projections_rebuilt",
"scratch_database": scratch,
"through_sequence": report.through_sequence.map(|seq| seq.value().to_string()),
"live_hash": report.live_hash,
"rebuilt_hash": report.rebuilt_hash,
"agrees": report.agrees,
}),
pretty,
);
if !report.agrees {
return Err(Failure::new(
KernelErrorCode::Storage,
"the rebuilt projections do not agree with the live ones",
));
}
Ok(())
}
pub async fn retention(what: &Retention, pretty: bool) -> Result<(), Failure> {
let config = AdminConfig::from_env().map_err(configuration)?;
let blob_config = BlobConfig::from_env().map_err(configuration)?;
let pool = gwk_kernel::connect_pool(config.admin_database_url(), 4)
.await
.map_err(configuration)?;
let blobs = PgBlobStore::open(pool, blob_config)
.await
.map_err(blob_failure)?;
let answer = match what {
Retention::Pin { address, evidence } => {
blobs
.pin(address, &EvidenceId::new(evidence.clone()))
.await
.map_err(blob_failure)?;
json!({"type": "blob_pinned", "address": address.as_str(), "evidence": evidence})
}
Retention::Unpin { address, evidence } => {
blobs
.unpin(address, &EvidenceId::new(evidence.clone()))
.await
.map_err(blob_failure)?;
json!({"type": "blob_unpinned", "address": address.as_str(), "evidence": evidence})
}
Retention::Sweep => {
let removed = blobs.sweep().await.map_err(blob_failure)?;
let addresses: Vec<&str> = removed.iter().map(BlobAddress::as_str).collect();
json!({"type": "blobs_swept", "removed": addresses})
}
Retention::Shred { address } => {
blobs.shred(address).await.map_err(blob_failure)?;
json!({"type": "blob_shredded", "address": address.as_str()})
}
};
emit(&answer, pretty);
Ok(())
}
fn revision() -> Result<String, Failure> {
if let Some(stamped) = PUBLIC_REVISION {
return Ok(stamped.to_owned());
}
let supplied = std::env::var(REVISION_ENV).ok().filter(|value| {
value.len() == 40
&& value
.bytes()
.all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
});
supplied.ok_or_else(|| {
Failure::usage(format!(
"this build carries no public revision, so it cannot record or report which build it \
is; rebuild from a clean checkout or set {REVISION_ENV} to a 40-character lowercase \
hexadecimal revision"
))
})
}
const REVISION_ENV: &str = "GWK_PUBLIC_REVISION";
fn beside(url: &SecretString, database: &str) -> Result<SecretString, Failure> {
if database.is_empty()
|| !database
.bytes()
.all(|b| b.is_ascii_alphanumeric() || b == b'_')
{
return Err(Failure::usage(format!(
"{database:?} is not a database name"
)));
}
let (prefix, tail) = url
.expose_secret()
.rsplit_once('/')
.ok_or_else(|| Failure::usage("the admin DSN has no /database to replace"))?;
let query = tail.find('?').map(|at| &tail[at..]).unwrap_or("");
Ok(SecretString::from(format!("{prefix}/{database}{query}")))
}
fn configuration(error: gwk_kernel::KernelError) -> Failure {
Failure::new(KernelErrorCode::Storage, error.to_string())
}
fn refusal(refusal: Refusal) -> Failure {
Failure::new(refusal.code, refusal.message)
}
fn blob_failure(error: gwk_domain::port::BlobError) -> Failure {
Failure::new(
match &error {
gwk_domain::port::BlobError::NotFound => KernelErrorCode::NotFound,
gwk_domain::port::BlobError::Tombstoned => KernelErrorCode::BlobTombstoned,
gwk_domain::port::BlobError::DigestMismatch { .. }
| gwk_domain::port::BlobError::Integrity(_) => KernelErrorCode::BlobIntegrity,
gwk_domain::port::BlobError::Pinned => KernelErrorCode::Authority,
gwk_domain::port::BlobError::Storage(_) => KernelErrorCode::Storage,
},
error.to_string(),
)
}
fn verdict(report: &recover::RecoveryReport) -> Value {
match &report.verdict {
recover::Verdict::Verified { anchor } => {
json!({"verdict": "verified", "anchor": anchor.value().to_string()})
}
recover::Verdict::Replayed { events } => {
json!({"verdict": "replayed", "events": events})
}
recover::Verdict::Unverified { reason } => {
json!({"verdict": "unverified", "reason": reason})
}
recover::Verdict::Diverged { expected, found } => {
json!({"verdict": "diverged", "expected": expected, "found": found})
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_scratch_database_lives_beside_the_one_it_verifies() {
let url = SecretString::from("postgres://gw@localhost:5432/gwk_live".to_owned());
let scratch = beside(&url, "gwk_scratch").expect("derive");
assert_eq!(
scratch.expose_secret(),
"postgres://gw@localhost:5432/gwk_scratch"
);
}
#[test]
fn the_scratch_connects_on_the_same_terms_as_the_live_one() {
let url = SecretString::from(
"postgres://gw@localhost:5432/gwk_live?sslmode=require&connect_timeout=5".to_owned(),
);
assert_eq!(
beside(&url, "gwk_scratch").expect("derive").expose_secret(),
"postgres://gw@localhost:5432/gwk_scratch?sslmode=require&connect_timeout=5"
);
}
#[test]
fn a_scratch_name_that_is_not_a_name_is_refused() {
let url = SecretString::from("postgres://gw@localhost:5432/gwk_live".to_owned());
for name in ["", "gwk scratch", "gwk;drop", "other/db", "gwk-scratch"] {
assert_eq!(
beside(&url, name).expect_err(name).exit,
crate::exit::USAGE,
"accepted {name:?}"
);
}
}
}