use gwk_domain::command::KernelCommand;
use gwk_domain::envelope::{
CommandEnvelope, ENVELOPE_SCHEMA_VERSION, EventEnvelope, accept_schema_version,
};
use gwk_domain::fsm::{AttemptState, CommandState, MessageState, TaskState};
use gwk_domain::ids::{AggregateId, EventId, ReceiptId, Seq};
use gwk_domain::port::EventStore;
use gwk_domain::protocol::{KernelErrorCode, KernelResult};
use gwk_domain::transition::{self, Cursor, TransitionRequest, TransitionResult};
use serde::de::DeserializeOwned;
use sqlx::{PgConnection, Row};
use crate::authority;
use crate::epoch::{self, Epoch};
use crate::numeric::from_numeric_text;
use crate::project::{
Refusal, apply_event, from_wire_str, page_attention, unresolved_attention, wire_str,
write_receipt,
};
use crate::store::{MAX_INFLIGHT_APPENDS, PgEventStore, current_aggregate_version, events_for_key};
const TASK_CURSOR: &str = "SELECT state, version FROM gwk.task WHERE id = $1";
const ATTEMPT_CURSOR: &str = "SELECT state, version FROM gwk.attempt WHERE id = $1";
const ATTEMPT_VERSION: &str = "SELECT version FROM gwk.attempt WHERE id = $1";
const LEASE_VERSION: &str = "SELECT version FROM gwk.lease WHERE id = $1";
const DISPATCH_NODE_VERSION: &str = "SELECT version FROM gwk.dispatch_node WHERE id = $1";
const AGGREGATE_OWNER: &str = "SELECT project_id FROM gwk.event \
WHERE aggregate_type = $1 AND aggregate_id = $2 \
ORDER BY aggregate_version LIMIT 1";
const MESSAGE_CURSOR: &str = "SELECT state, version FROM gwk.message WHERE id = $1";
const COMMAND_CURSOR: &str = "SELECT state, version FROM gwk.command WHERE id = $1";
const GATE_VERSION: &str = "SELECT version FROM gwk.gate WHERE id = $1";
const CHECKPOINT_SEQ: &str =
"SELECT seq::text AS seq_text FROM gwk.orchestrator_checkpoint WHERE orchestrator_id = $1";
#[derive(Debug)]
struct Route {
aggregate_type: &'static str,
aggregate_id: String,
event_type: &'static str,
}
enum Prior {
Unused,
Replay(Vec<EventEnvelope>),
Conflict(String),
}
impl PgEventStore {
pub async fn submit(&self, envelope: &CommandEnvelope) -> KernelResult {
match self.try_submit(envelope).await {
Ok(result) => result,
Err(refusal) => refusal.into_result(),
}
}
async fn try_submit(&self, envelope: &CommandEnvelope) -> Result<KernelResult, Refusal> {
let _permit = self.admit().map_err(|_| {
Refusal::new(
KernelErrorCode::Overloaded,
format!("append queue is full ({MAX_INFLIGHT_APPENDS} in flight)"),
)
})?;
accept_schema_version(envelope.schema_version, ENVELOPE_SCHEMA_VERSION, &[])
.map_err(|e| Refusal::new(KernelErrorCode::Schema, e.to_string()))?;
let command = KernelCommand::from_envelope(envelope)
.map_err(|e| Refusal::validation(e.to_string()))?;
let route = route_of(envelope, &command)?;
check_routing(envelope, &route)?;
check_body_project(envelope, &command)?;
check_activation(envelope, &command)?;
let payload = serde_json::to_value(&command)
.map_err(|e| Refusal::storage(format!("serialize command body: {e}")))?;
let mut tx = self
.pool()
.begin()
.await
.map_err(|e| Refusal::storage(format!("begin: {e}")))?;
let writer = self.lock_writer(&mut tx).await?;
let epoch = epoch::epoch_of(&mut tx).await?;
if !admitted(epoch, &command) {
return Err(epoch::sealed_refusal(epoch, command.command_type()));
}
let expected_version = match prior_for_key(&mut tx, envelope, &route, &payload).await? {
Prior::Replay(events) => {
tx.commit()
.await
.map_err(|e| Refusal::storage(format!("commit: {e}")))?;
return self.applied(envelope, events).await;
}
Prior::Conflict(reason) => {
return Err(Refusal::new(KernelErrorCode::IdempotencyConflict, reason));
}
Prior::Unused => {
check_aggregate_owner(&mut tx, envelope, &route).await?;
check_second_cutover(&mut tx, &command, epoch).await?;
let decision = authority::evaluate(
&mut tx,
envelope,
&envelope.command_type,
&route.aggregate_id,
)
.await?;
if let authority::Decision::Page { action_class } = decision {
return self.paged(tx, envelope, &route, action_class).await;
}
if let Some(action_class) = decision.action_class() {
write_receipt(
&mut tx,
&receipt_for(
envelope,
&route,
action_class,
"matching unexpired scoped grant",
),
)
.await?;
}
check_attention_dedup(&mut tx, &command).await?;
decide(&mut tx, envelope, &command, &route).await?
}
};
check_expected_version(envelope, expected_version)?;
let event = build_event(envelope, payload, &route, expected_version)?;
let fence = writer.current_fence.map(gwk_domain::ids::FenceToken::new);
let appended = self
.append_locked(&mut tx, &writer, expected_version, fence, &[event])
.await?;
if !appended.replayed {
for event in &appended.events {
apply_event(&mut tx, event).await?;
}
}
if let Some(last) = appended.events.last() {
self.checkpoint_if_due(&mut tx, &writer, last.global_sequence, &last.appended_at)
.await?;
}
tx.commit()
.await
.map_err(|e| Refusal::storage(format!("commit: {e}")))?;
self.applied(envelope, appended.events).await
}
async fn paged(
&self,
mut tx: sqlx::Transaction<'_, sqlx::Postgres>,
envelope: &CommandEnvelope,
route: &Route,
action_class: &'static str,
) -> Result<KernelResult, Refusal> {
let actor = actor_json(envelope)?;
let at = envelope.issued_at.as_str();
let subject = format!("{}/{}", route.aggregate_type, route.aggregate_id);
write_receipt(
&mut tx,
&receipt_for(
envelope,
route,
action_class,
"no matching unexpired scoped grant",
),
)
.await?;
page_attention(
&mut tx,
&format!("page:{action_class}:{subject}"),
&format!(
"{} requires an unexpired {action_class} grant",
envelope.command_type
),
&subject,
&actor,
at,
)
.await?;
tx.commit()
.await
.map_err(|e| Refusal::storage(format!("commit: {e}")))?;
Err(Refusal::new(
KernelErrorCode::Authority,
format!(
"{} on {subject} requires an unexpired {action_class} grant; attention raised",
envelope.command_type
),
)
.with_detail(serde_json::json!({ "action_class": action_class, "subject": subject })))
}
async fn applied(
&self,
envelope: &CommandEnvelope,
events: Vec<EventEnvelope>,
) -> Result<KernelResult, Refusal> {
let watermark = self.watermark().await?.ok_or_else(|| {
Refusal::storage("the log is empty after a committed append".to_owned())
})?;
Ok(KernelResult::CommandApplied {
command_id: envelope.command_id.clone(),
events,
watermark,
})
}
}
fn route_of(envelope: &CommandEnvelope, command: &KernelCommand) -> Result<Route, Refusal> {
use KernelCommand as C;
let (aggregate_type, aggregate_id, event_type) = match command {
C::ActivateKernel { .. } => (
epoch::KERNEL_AGGREGATE,
epoch::KERNEL_SINGLETON.to_owned(),
epoch::ACTIVATION_EVENT_TYPE,
),
C::CreateTask { task_id, .. } => ("task", task_id.as_str().to_owned(), "task_created"),
C::TransitionTask { task_id, .. } => {
("task", task_id.as_str().to_owned(), "task_transitioned")
}
C::CreateAttempt { attempt_id, .. } => {
("attempt", attempt_id.as_str().to_owned(), "attempt_created")
}
C::TransitionAttempt { attempt_id, .. } => (
"attempt",
attempt_id.as_str().to_owned(),
"attempt_transitioned",
),
C::RecordAttemptOutcome { attempt_id, .. } => (
"attempt",
attempt_id.as_str().to_owned(),
"attempt_outcome_recorded",
),
C::UpdateBudget { attempt_id, .. } => {
("attempt", attempt_id.as_str().to_owned(), "budget_updated")
}
C::RecordRound { attempt_id, .. } => {
("attempt", attempt_id.as_str().to_owned(), "round_recorded")
}
C::RecordFinding { attempt_id, .. } => (
"attempt",
attempt_id.as_str().to_owned(),
"finding_recorded",
),
C::OpenEngineSession {
engine_session_id, ..
} => (
"engine_session",
engine_session_id.as_str().to_owned(),
"engine_session_opened",
),
C::CloseEngineSession { engine_session_id } => (
"engine_session",
engine_session_id.as_str().to_owned(),
"engine_session_closed",
),
C::AcquireLease { lease_id, .. } => {
("lease", lease_id.as_str().to_owned(), "lease_acquired")
}
C::RenewLease { lease_id, .. } => ("lease", lease_id.as_str().to_owned(), "lease_renewed"),
C::ReleaseLease { lease_id, .. } => {
("lease", lease_id.as_str().to_owned(), "lease_released")
}
C::ExpireLease { lease_id, .. } => ("lease", lease_id.as_str().to_owned(), "lease_expired"),
C::RegisterWorktree { worktree_id, .. } => (
"worktree",
worktree_id.as_str().to_owned(),
"worktree_registered",
),
C::UpdateWorktree { worktree_id, .. } => (
"worktree",
worktree_id.as_str().to_owned(),
"worktree_updated",
),
C::ReleaseWorktree { worktree_id, .. } => (
"worktree",
worktree_id.as_str().to_owned(),
"worktree_released",
),
C::RegisterDispatchNode {
dispatch_node_id, ..
} => (
"dispatch_node",
dispatch_node_id.as_str().to_owned(),
"dispatch_node_registered",
),
C::TransitionDispatchNode {
dispatch_node_id, ..
} => (
"dispatch_node",
dispatch_node_id.as_str().to_owned(),
"dispatch_node_transitioned",
),
C::SendMessage { message_id, .. } => {
("message", message_id.as_str().to_owned(), "message_sent")
}
C::TransitionMessage { message_id, .. } => (
"message",
message_id.as_str().to_owned(),
"message_transitioned",
),
C::IssueCommand { command_id, .. } => {
("command", command_id.as_str().to_owned(), "command_issued")
}
C::TransitionCommand { command_id, .. } => (
"command",
command_id.as_str().to_owned(),
"command_transitioned",
),
C::RecordCommandOutcome { command_id, .. } => (
"command",
command_id.as_str().to_owned(),
"command_outcome_recorded",
),
C::OpenGate { gate_id, .. } => ("gate", gate_id.as_str().to_owned(), "gate_opened"),
C::DecideGate { gate_id, .. } => ("gate", gate_id.as_str().to_owned(), "gate_decided"),
C::RecordEvidence { evidence_id, .. } => (
"evidence",
evidence_id.as_str().to_owned(),
"evidence_recorded",
),
C::GrantAuthority {
authority_grant_id, ..
} => (
"authority_grant",
authority_grant_id.as_str().to_owned(),
"authority_granted",
),
C::RevokeAuthority {
authority_grant_id, ..
} => (
"authority_grant",
authority_grant_id.as_str().to_owned(),
"authority_revoked",
),
C::RaiseAttention {
attention_item_id, ..
} => (
"attention_item",
attention_item_id.as_str().to_owned(),
"attention_raised",
),
C::ResolveAttention {
attention_item_id, ..
} => (
"attention_item",
attention_item_id.as_str().to_owned(),
"attention_resolved",
),
C::WriteOrchestratorCheckpoint { checkpoint } => (
"orchestrator_checkpoint",
checkpoint.orchestrator_id.clone().ok_or_else(|| {
Refusal::validation("a checkpoint without an orchestrator_id has no identity")
})?,
"orchestrator_checkpoint_written",
),
C::IngestRecord { .. } => (
"ingested_record",
format!(
"ingest:{}:{}",
envelope.project_id.as_str(),
envelope.idempotency_key.as_str()
),
"record_ingested",
),
};
if aggregate_id.is_empty() {
return Err(Refusal::validation(format!(
"{} names an empty {aggregate_type} id",
command.command_type()
)));
}
Ok(Route {
aggregate_type,
aggregate_id,
event_type,
})
}
fn check_routing(envelope: &CommandEnvelope, route: &Route) -> Result<(), Refusal> {
let mismatch = |field: &str, declared: &str| {
Refusal::validation(format!(
"envelope {field} {declared:?} does not name the aggregate the body addresses \
({}/{})",
route.aggregate_type, route.aggregate_id
))
};
if let Some(declared) = &envelope.target_aggregate_type
&& declared != route.aggregate_type
{
return Err(mismatch("target_aggregate_type", declared));
}
if let Some(declared) = &envelope.target_aggregate_id
&& declared.as_str() != route.aggregate_id
{
return Err(mismatch("target_aggregate_id", declared.as_str()));
}
Ok(())
}
fn receipt_for(
envelope: &CommandEnvelope,
route: &Route,
action_class: &str,
observed_basis: &str,
) -> gwk_domain::entity::Receipt {
gwk_domain::entity::Receipt {
id: ReceiptId::new(format!(
"receipt:{}:{}",
envelope.project_id.as_str(),
envelope.idempotency_key.as_str()
)),
actor: envelope.actor.clone(),
action: action_class.to_owned(),
subject_type: route.aggregate_type.to_owned(),
subject_id: route.aggregate_id.clone(),
from: None,
to: None,
observed_basis: Some(observed_basis.to_owned()),
ts: envelope.issued_at.clone(),
}
}
fn actor_json(envelope: &CommandEnvelope) -> Result<serde_json::Value, Refusal> {
serde_json::to_value(&envelope.actor)
.map_err(|e| Refusal::storage(format!("serialize actor: {e}")))
}
async fn check_attention_dedup(
conn: &mut PgConnection,
command: &KernelCommand,
) -> Result<(), Refusal> {
let KernelCommand::RaiseAttention {
kind, subject_ref, ..
} = command
else {
return Ok(());
};
if let Some(open) = unresolved_attention(conn, kind, subject_ref.as_deref()).await? {
return Err(Refusal::validation(format!(
"attention item {open:?} is already open on ({kind:?}, {:?})",
subject_ref.as_deref().unwrap_or_default()
)));
}
Ok(())
}
fn check_body_project(envelope: &CommandEnvelope, command: &KernelCommand) -> Result<(), Refusal> {
let KernelCommand::CreateTask {
project: Some(project),
..
} = command
else {
return Ok(());
};
if project != envelope.project_id.as_str() {
return Err(Refusal::validation(format!(
"body project {project:?} does not name the envelope's project {:?}",
envelope.project_id.as_str()
)));
}
Ok(())
}
fn admitted(epoch: Epoch, command: &KernelCommand) -> bool {
match epoch {
Epoch::Active => true,
Epoch::Sealed => matches!(command, KernelCommand::ActivateKernel { .. }),
Epoch::None => false,
}
}
fn check_activation(envelope: &CommandEnvelope, command: &KernelCommand) -> Result<(), Refusal> {
let KernelCommand::ActivateKernel {
cutover_id,
archive_manifest_sha256,
} = command
else {
return Ok(());
};
if cutover_id.is_empty() {
return Err(Refusal::validation(
"an activation with an empty cutover id names no cutover",
));
}
if !gwk_domain::blob::is_sha256_hex(archive_manifest_sha256) {
return Err(Refusal::validation(format!(
"archive_manifest_sha256 must be a lowercase 64-hex digest, not \
{archive_manifest_sha256:?}"
)));
}
let required = epoch::activation_key(cutover_id);
if envelope.idempotency_key.as_str() != required {
return Err(Refusal::validation(format!(
"activating cutover {cutover_id:?} requires idempotency key {required:?}, not {:?}",
envelope.idempotency_key.as_str()
)));
}
Ok(())
}
async fn check_second_cutover(
conn: &mut PgConnection,
command: &KernelCommand,
epoch: Epoch,
) -> Result<(), Refusal> {
let KernelCommand::ActivateKernel { cutover_id, .. } = command else {
return Ok(());
};
if epoch != Epoch::Active {
return Ok(());
}
let committed = epoch::committed_cutover(conn).await?.ok_or_else(|| {
Refusal::storage("the kernel is past genesis with no activation event".to_owned())
})?;
Err(Refusal::new(
KernelErrorCode::AlreadyActive,
format!(
"this kernel activated at cutover {committed:?}; {cutover_id:?} is a different \
cutover"
),
)
.with_detail(serde_json::json!({ "activated_cutover_id": committed })))
}
async fn check_aggregate_owner(
conn: &mut PgConnection,
envelope: &CommandEnvelope,
route: &Route,
) -> Result<(), Refusal> {
let owner: Option<String> = sqlx::query_scalar(AGGREGATE_OWNER)
.bind(route.aggregate_type)
.bind(&route.aggregate_id)
.fetch_optional(conn)
.await
.map_err(|e| Refusal::storage(format!("read aggregate owner: {e}")))?;
match owner {
Some(owner) if owner != envelope.project_id.as_str() => Err(Refusal::validation(format!(
"{}/{} belongs to project {owner:?}, not {:?}",
route.aggregate_type,
route.aggregate_id,
envelope.project_id.as_str()
))),
_ => Ok(()),
}
}
fn check_expected_version(envelope: &CommandEnvelope, derived: u32) -> Result<(), Refusal> {
match envelope.expected_version {
Some(expected) if expected != derived => Err(Refusal::new(
KernelErrorCode::StaleVersion,
format!("version conflict: actual {derived}, expected {expected}"),
)
.with_detail(serde_json::json!({ "actual": derived, "expected": expected }))),
_ => Ok(()),
}
}
async fn prior_for_key(
conn: &mut PgConnection,
envelope: &CommandEnvelope,
route: &Route,
payload: &serde_json::Value,
) -> Result<Prior, Refusal> {
let stored = events_for_key(
conn,
envelope.project_id.as_str(),
route.aggregate_type,
&route.aggregate_id,
envelope.idempotency_key.as_str(),
)
.await?;
let Some(first) = stored.first() else {
return Ok(Prior::Unused);
};
let same_request = stored.len() == 1
&& first.project_id == envelope.project_id
&& first.aggregate_type == route.aggregate_type
&& first.aggregate_id.as_str() == route.aggregate_id
&& first.actor == envelope.actor
&& &first.payload == payload;
if same_request {
return Ok(Prior::Replay(stored));
}
Ok(Prior::Conflict(format!(
"idempotency key {:?} already names a different request on {}/{} in project {}",
envelope.idempotency_key.as_str(),
first.aggregate_type,
first.aggregate_id,
first.project_id
)))
}
async fn decide(
conn: &mut PgConnection,
envelope: &CommandEnvelope,
command: &KernelCommand,
route: &Route,
) -> Result<u32, Refusal> {
use KernelCommand as C;
Ok(match command {
C::ActivateKernel { .. } => 1,
C::CreateTask { .. }
| C::CreateAttempt { .. }
| C::OpenEngineSession { .. }
| C::AcquireLease { .. }
| C::RegisterWorktree { .. }
| C::RegisterDispatchNode { .. }
| C::SendMessage { .. }
| C::IssueCommand { .. }
| C::OpenGate { .. }
| C::RecordEvidence { .. }
| C::GrantAuthority { .. }
| C::RaiseAttention { .. } => 0,
C::RevokeAuthority { .. } | C::ResolveAttention { .. } => {
current_aggregate_version(conn, route.aggregate_type, &route.aggregate_id).await?
}
C::TransitionTask {
to,
expected_version,
..
} => {
let cursor: Cursor<TaskState> =
fsm_cursor(conn, TASK_CURSOR, &route.aggregate_id, "task").await?;
decide_transition(&cursor, *to, *expected_version, envelope, None)?
}
C::TransitionAttempt {
to,
expected_version,
receipt_id,
..
} => {
let cursor: Cursor<AttemptState> =
fsm_cursor(conn, ATTEMPT_CURSOR, &route.aggregate_id, "attempt").await?;
decide_transition(
&cursor,
*to,
*expected_version,
envelope,
receipt_id.as_ref(),
)?
}
C::TransitionMessage {
to,
expected_version,
..
} => {
let cursor: Cursor<MessageState> =
fsm_cursor(conn, MESSAGE_CURSOR, &route.aggregate_id, "message").await?;
decide_transition(&cursor, *to, *expected_version, envelope, None)?
}
C::TransitionCommand {
to,
expected_version,
..
} => {
if *to == CommandState::VerificationComplete {
return Err(Refusal::validation(
"a command reaches verification_complete by recording its outcome, \
which is the same write",
));
}
let cursor: Cursor<CommandState> =
fsm_cursor(conn, COMMAND_CURSOR, &route.aggregate_id, "command").await?;
decide_transition(&cursor, *to, *expected_version, envelope, None)?
}
C::RecordCommandOutcome {
expected_version, ..
} => {
let cursor: Cursor<CommandState> =
fsm_cursor(conn, COMMAND_CURSOR, &route.aggregate_id, "command").await?;
decide_transition(
&cursor,
CommandState::VerificationComplete,
*expected_version,
envelope,
None,
)?
}
C::DecideGate {
expected_version, ..
} => {
decide_cas(
conn,
GATE_VERSION,
&route.aggregate_id,
"gate",
*expected_version,
)
.await?
}
C::UpdateBudget {
expected_version, ..
}
| C::RecordAttemptOutcome {
expected_version, ..
} => {
decide_cas(
conn,
ATTEMPT_VERSION,
&route.aggregate_id,
"attempt",
*expected_version,
)
.await?
}
C::RenewLease {
expected_version, ..
}
| C::ReleaseLease {
expected_version, ..
}
| C::ExpireLease {
expected_version, ..
} => {
decide_cas(
conn,
LEASE_VERSION,
&route.aggregate_id,
"lease",
*expected_version,
)
.await?
}
C::TransitionDispatchNode {
expected_version, ..
} => {
decide_cas(
conn,
DISPATCH_NODE_VERSION,
&route.aggregate_id,
"dispatch_node",
*expected_version,
)
.await?
}
C::WriteOrchestratorCheckpoint { checkpoint } => {
if let Some(current) = checkpoint_seq(conn, &route.aggregate_id).await?
&& checkpoint.seq.value() <= current
{
return Err(Refusal::new(
KernelErrorCode::StaleVersion,
format!(
"checkpoint seq must advance past {current} (got {})",
checkpoint.seq
),
)
.with_detail(serde_json::json!({
"actual": current.to_string(),
"presented": checkpoint.seq.to_string(),
})));
}
current_aggregate_version(conn, route.aggregate_type, &route.aggregate_id).await?
}
C::CloseEngineSession { .. }
| C::UpdateWorktree { .. }
| C::ReleaseWorktree { .. }
| C::RecordRound { .. }
| C::RecordFinding { .. } => {
current_aggregate_version(conn, route.aggregate_type, &route.aggregate_id).await?
}
C::IngestRecord { .. } => {
current_aggregate_version(conn, route.aggregate_type, &route.aggregate_id).await?
}
})
}
async fn fsm_cursor<S: DeserializeOwned>(
conn: &mut PgConnection,
select: &'static str,
id: &str,
kind: &str,
) -> Result<Cursor<S>, Refusal> {
let row = sqlx::query(select)
.bind(id)
.fetch_optional(conn)
.await
.map_err(|e| Refusal::storage(format!("read {kind} {id}: {e}")))?
.ok_or_else(|| Refusal::not_found(format!("no {kind} {id}")))?;
let state: String = row
.try_get("state")
.map_err(|e| Refusal::storage(format!("column state: {e}")))?;
Ok(Cursor {
state: from_wire_str(&state)?,
version: row_version(&row)?,
applied_idempotency_key: None,
applied_by: None,
})
}
fn row_version(row: &sqlx::postgres::PgRow) -> Result<u32, Refusal> {
let version: i64 = row
.try_get("version")
.map_err(|e| Refusal::storage(format!("column version: {e}")))?;
u32::try_from(version).map_err(|e| Refusal::storage(format!("version out of range: {e}")))
}
fn decide_transition<S>(
cursor: &Cursor<S>,
to: S,
expected_version: u32,
envelope: &CommandEnvelope,
receipt_id: Option<&ReceiptId>,
) -> Result<u32, Refusal>
where
S: transition::TransitionGuard + serde::Serialize,
{
let request = TransitionRequest {
to,
expected_version,
actor: &envelope.actor,
idempotency_key: None,
receipt_id,
};
match transition::apply(cursor, &request) {
TransitionResult::Applied { .. } => Ok(cursor.version),
TransitionResult::IllegalEdge { from, to } => {
let (from, to) = (wire_str(&from)?, wire_str(&to)?);
Err(
Refusal::new(KernelErrorCode::IllegalEdge, format!("{from} -> {to}"))
.with_detail(serde_json::json!({ "from": from, "to": to })),
)
}
TransitionResult::StaleVersion { actual, expected } => Err(Refusal::new(
KernelErrorCode::StaleVersion,
format!("version conflict: actual {actual}, expected {expected}"),
)
.with_detail(serde_json::json!({ "actual": actual, "expected": expected }))),
TransitionResult::UnauthorizedActor { reason } => {
Err(Refusal::new(KernelErrorCode::Authority, reason))
}
}
}
async fn decide_cas(
conn: &mut PgConnection,
select: &'static str,
id: &str,
kind: &str,
expected: u32,
) -> Result<u32, Refusal> {
let row = sqlx::query(select)
.bind(id)
.fetch_optional(conn)
.await
.map_err(|e| Refusal::storage(format!("read {kind} {id}: {e}")))?
.ok_or_else(|| Refusal::not_found(format!("no {kind} {id}")))?;
let actual = row_version(&row)?;
if actual != expected {
return Err(Refusal::new(
KernelErrorCode::StaleVersion,
format!("version conflict: actual {actual}, expected {expected}"),
)
.with_detail(serde_json::json!({ "actual": actual, "expected": expected })));
}
Ok(actual)
}
async fn checkpoint_seq(
conn: &mut PgConnection,
orchestrator_id: &str,
) -> Result<Option<u64>, Refusal> {
let text: Option<String> = sqlx::query_scalar(CHECKPOINT_SEQ)
.bind(orchestrator_id)
.fetch_optional(conn)
.await
.map_err(|e| Refusal::storage(format!("read checkpoint {orchestrator_id}: {e}")))?;
text.map(|text| from_numeric_text(&text))
.transpose()
.map_err(|e| Refusal::storage(format!("column seq: {e}")))
}
fn build_event(
envelope: &CommandEnvelope,
payload: serde_json::Value,
route: &Route,
expected_version: u32,
) -> Result<EventEnvelope, Refusal> {
let aggregate_version = expected_version.checked_add(1).ok_or_else(|| {
Refusal::new(
KernelErrorCode::StaleVersion,
format!(
"{}/{} is at the version ceiling",
route.aggregate_type, route.aggregate_id
),
)
})?;
Ok(EventEnvelope {
event_id: EventId::new(format!(
"{}:{}:{aggregate_version}",
route.aggregate_type, route.aggregate_id
)),
project_id: envelope.project_id.clone(),
aggregate_type: route.aggregate_type.to_owned(),
aggregate_id: AggregateId::new(route.aggregate_id.clone()),
aggregate_version,
event_type: route.event_type.to_owned(),
schema_version: ENVELOPE_SCHEMA_VERSION,
global_sequence: Seq::new(0),
appended_at: envelope.issued_at.clone(),
occurred_at: envelope.issued_at.clone(),
actor: envelope.actor.clone(),
origin: envelope.origin.clone(),
causation_id: envelope.causation_id.clone(),
correlation_id: envelope.correlation_id.clone(),
idempotency_key: Some(envelope.idempotency_key.clone()),
payload,
payload_ref: None,
})
}
#[cfg(test)]
mod tests {
use gwk_domain::ids::{AttemptId, TaskId};
use super::*;
fn envelope(command: &KernelCommand) -> CommandEnvelope {
CommandEnvelope {
command_id: gwk_domain::ids::CommandId::new("cmd-1"),
project_id: gwk_domain::ids::ProjectId::new("p"),
command_type: command.command_type().to_owned(),
schema_version: ENVELOPE_SCHEMA_VERSION,
issued_at: gwk_domain::ids::Timestamp::new("2026-07-28T00:00:00Z"),
actor: gwk_domain::envelope::Actor {
kind: "kernel".into(),
id: None,
},
origin: gwk_domain::envelope::Origin {
system: "gw".into(),
r#ref: None,
},
target_aggregate_type: None,
target_aggregate_id: None,
expected_version: None,
idempotency_key: gwk_domain::ids::IdempotencyKey::new("k-1"),
causation_id: None,
correlation_id: None,
payload: serde_json::json!({}),
}
}
#[test]
fn every_in_scope_command_routes_to_an_aggregate_and_an_event_name() {
let create = KernelCommand::CreateTask {
task_id: TaskId::new("t-1"),
kind: None,
title: None,
spec_ref: None,
project: None,
priority: None,
tracker_ref: None,
};
let route = route_of(&envelope(&create), &create).expect("routes");
assert_eq!(route.aggregate_type, "task");
assert_eq!(route.aggregate_id, "t-1");
assert_eq!(route.event_type, "task_created");
let round = KernelCommand::RecordRound {
attempt_id: AttemptId::new("a-1"),
round: 2,
findings: gwk_domain::inherited::RoundFindingSummary {
total: 0,
auto_fix: 0,
ask_user: 0,
no_op: 0,
},
};
let route = route_of(&envelope(&round), &round).expect("routes");
assert_eq!(
(route.aggregate_type, route.aggregate_id.as_str()),
("attempt", "a-1")
);
let evidence = KernelCommand::RecordEvidence {
evidence_id: gwk_domain::ids::EvidenceId::new("ev-1"),
kind: "diff".to_owned(),
r#ref: "blob://d".to_owned(),
digest: None,
byte_size: None,
};
let route = route_of(&envelope(&evidence), &evidence).expect("routes");
assert_eq!(
(route.aggregate_type, route.event_type),
("evidence", "evidence_recorded")
);
}
#[test]
fn an_ingested_record_is_named_by_the_envelope_that_carried_it() {
let ingest = KernelCommand::IngestRecord {
kind: gwk_domain::ingestion::IngestionKind::Memory,
payload: serde_json::json!({ "text": "recalled" }),
payload_ref: None,
};
let route = route_of(&envelope(&ingest), &ingest).expect("routes");
assert_eq!(
(
route.aggregate_type,
route.aggregate_id.as_str(),
route.event_type
),
("ingested_record", "ingest:p:k-1", "record_ingested")
);
let mut second = envelope(&ingest);
second.idempotency_key = gwk_domain::ids::IdempotencyKey::new("k-2");
assert_eq!(
route_of(&second, &ingest).expect("routes").aggregate_id,
"ingest:p:k-2"
);
}
#[test]
fn an_anonymous_checkpoint_has_no_row_to_be_the_latest_of() {
let anonymous = KernelCommand::WriteOrchestratorCheckpoint {
checkpoint: gwk_domain::inherited::OrchestratorCheckpoint {
orchestrator_id: None,
seq: Seq::new(1),
native_session_ref: None,
active_goal: None,
active_step_ref: None,
latest_command_ref: None,
open_attempts: None,
leases: None,
pending_approvals: None,
budget_cursor: None,
},
};
let refusal = route_of(&envelope(&anonymous), &anonymous).expect_err("no identity");
assert_eq!(refusal.code, KernelErrorCode::Validation);
}
#[test]
fn an_empty_id_is_refused_before_it_becomes_an_aggregate() {
let nameless = KernelCommand::CreateTask {
task_id: TaskId::new(""),
kind: None,
title: None,
spec_ref: None,
project: None,
priority: None,
tracker_ref: None,
};
assert_eq!(
route_of(&envelope(&nameless), &nameless)
.expect_err("empty id")
.code,
KernelErrorCode::Validation
);
}
#[test]
fn the_event_carries_the_kernels_schema_version_and_a_derived_id() {
let command = KernelCommand::TransitionTask {
task_id: TaskId::new("t-1"),
to: TaskState::Working,
expected_version: 1,
};
let route = route_of(&envelope(&command), &command).expect("routes");
let payload = serde_json::to_value(&command).expect("serialize");
let event = build_event(&envelope(&command), payload.clone(), &route, 1).expect("built");
assert_eq!(event.aggregate_version, 2);
assert_eq!(event.event_id.as_str(), "task:t-1:2");
assert_eq!(event.schema_version, ENVELOPE_SCHEMA_VERSION);
assert_eq!(event.event_type, "task_transitioned");
assert_eq!(event.payload, payload);
assert_eq!(
event.idempotency_key.as_ref().map(|k| k.as_str()),
Some("k-1")
);
}
#[test]
fn a_transition_at_the_version_ceiling_is_refused_not_wrapped() {
let command = KernelCommand::TransitionTask {
task_id: TaskId::new("t-1"),
to: TaskState::Working,
expected_version: u32::MAX,
};
let route = route_of(&envelope(&command), &command).expect("routes");
let refusal = build_event(&envelope(&command), serde_json::json!({}), &route, u32::MAX)
.expect_err("the ceiling");
assert_eq!(refusal.code, KernelErrorCode::StaleVersion);
}
#[test]
fn the_liveness_flip_refusal_is_authority_not_a_version_probe() {
let command = KernelCommand::TransitionAttempt {
attempt_id: AttemptId::new("a-1"),
to: AttemptState::Blocked,
expected_version: 999,
receipt_id: None,
};
let cursor = Cursor {
state: AttemptState::Running,
version: 5,
applied_idempotency_key: None,
applied_by: None,
};
let refusal = decide_transition(
&cursor,
AttemptState::Blocked,
999,
&envelope(&command),
None,
)
.expect_err("wrong actor");
assert_eq!(refusal.code, KernelErrorCode::Authority);
assert!(!refusal.message.contains('5'), "{refusal}");
}
#[test]
fn a_refusal_carries_the_number_a_retrier_needs() {
let command = KernelCommand::TransitionTask {
task_id: TaskId::new("t-1"),
to: TaskState::Working,
expected_version: 1,
};
let cursor = Cursor {
state: TaskState::Submitted,
version: 4,
applied_idempotency_key: None,
applied_by: None,
};
let refusal = decide_transition(&cursor, TaskState::Working, 1, &envelope(&command), None)
.expect_err("stale");
assert_eq!(refusal.code, KernelErrorCode::StaleVersion);
assert_eq!(
refusal.detail,
Some(serde_json::json!({ "actual": 4, "expected": 1 }))
);
}
}