Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 0 additions & 1 deletion .github/workflows/guard.yml
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,6 @@ concurrency:
jobs:
repo:
name: repo
if: github.event_name != 'pull_request'
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v6
Expand Down
13 changes: 9 additions & 4 deletions .runseal/lib/guard/guard.ts
Original file line number Diff line number Diff line change
@@ -1,9 +1,9 @@
//! `runseal :guard` — the local validation suite the pre-commit hook runs.
//!
//! Mirrors the CI `smoke (ubuntu-latest)` job (cargo fmt/clippy/test, all
//! `--locked` like CI) so "green locally" means "green in CI", plus the Deno
//! checks for the `.runseal/` wrappers that CI does not cover. Runnable on
//! demand as well as from the hook.
//! Mirrors the CI repo + Rust checks (flavor and cargo, all `--locked` like CI)
//! so "green locally" means "green in CI", plus the Deno checks for the
//! `.runseal/` wrappers that CI does not cover. Runnable on demand as well as
//! from the hook.

import { run } from "@/lib/std/cmd.ts";
import { join } from "@/lib/std/fs.ts";
Expand All @@ -21,6 +21,11 @@ export async function guard(): Promise<number> {
const config = ".runseal/deno.json";

const steps: Step[] = [
{
title: "flavor check",
command: "flavor",
args: ["check", "--root", ".", "--config", "flavor.toml"],
},
{ title: "cargo fmt", command: "cargo", args: ["fmt", "--all", "--check"] },
{
title: "cargo clippy",
Expand Down
10 changes: 5 additions & 5 deletions .runseal/lib/land/land.ts
Original file line number Diff line number Diff line change
Expand Up @@ -72,10 +72,10 @@ export async function land(argv: string[]): Promise<number> {
}

// 4. Wait for required checks — polled here, quietly, with a bounded timeout.
console.log(`waiting for required checks (timeout ${CHECK_TIMEOUT_MS / 1000}s) ...`);
console.log(`waiting for PR checks (timeout ${CHECK_TIMEOUT_MS / 1000}s) ...`);
const outcome = await waitChecks(pr, repo);
if (outcome !== "pass") {
return fail(`required checks ${outcome} — PR #${pr} left open`);
return fail(`PR checks ${outcome} — PR #${pr} left open`);
}

// 5. Squash merge.
Expand Down Expand Up @@ -126,8 +126,8 @@ async function openPr(branch: string, repo: string): Promise<number | null> {
}

/**
* Poll the PR's required checks until they all pass, one fails, or the timeout
* elapses. Parses `--json bucket` from stdout and ignores exit codes so pending
* Poll all PR checks until they pass, one fails, or the timeout elapses. Parses
* `--json bucket` from stdout and ignores exit codes so pending
* states are handled uniformly; an empty set means checks have not registered
* yet (just after creation).
*/
Expand All @@ -136,7 +136,7 @@ async function waitChecks(pr: number, repo: string): Promise<"pass" | "fail" | "
while (Date.now() < deadline) {
const result = await capture(
"gh",
["pr", "checks", String(pr), "--required", "--json", "bucket"],
["pr", "checks", String(pr), "--json", "bucket"],
{ cwd: repo },
);
let buckets: string[] = [];
Expand Down
4 changes: 2 additions & 2 deletions .runseal/wrappers/guard.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
//! `runseal :guard` — run the local validation suite (mirrors CI smoke + Deno
//! checks). Thin entry point; logic lives in the guard module.
//! `runseal :guard` — run the local validation suite (mirrors CI repo/Rust +
//! Deno checks). Thin entry point; logic lives in the guard module.

import { guard } from "@/lib/guard/guard.ts";

Expand Down
1 change: 1 addition & 0 deletions crates/santi-api/src/server.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
mod error;
mod errors;
mod im;
mod ingress;
mod openapi;
mod routes;
Expand Down
64 changes: 64 additions & 0 deletions crates/santi-api/src/server/im.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
use axum::{
Json,
extract::{Path, Query, State},
};
use santi_core::{
ImInboxEntry, ImSendRequest, ImSendResponse, IngestOutcome, SantiError, SantiService,
};

use super::ApiError;

#[utoipa::path(
post,
path = "/api/v1/im/send",
request_body = ImSendRequest,
responses(
(status = 200, body = ImSendResponse),
(status = 423, body = SantiError),
(status = 404, body = SantiError),
(status = 500, body = SantiError)
)
)]
pub(super) async fn send_im(
State(service): State<SantiService>,
Json(request): Json<ImSendRequest>,
) -> Result<Json<ImSendResponse>, ApiError> {
let outcome = service
.im_send(&request.soul_id, &request.participant_id, &request.content)
.map_err(ApiError::from_service)?;
match outcome {
IngestOutcome::Accepted { receipt } => Ok(Json(ImSendResponse {
participant_id: request.participant_id,
receipt,
})),
IngestOutcome::Rejected { error } => Err(ApiError::from_santi(*error)),
}
}

#[utoipa::path(
get,
path = "/api/v1/im/inbox/{participant_id}",
params(
("participant_id" = String, Path),
("since" = Option<i64>, Query)
),
responses(
(status = 200, body = Vec<ImInboxEntry>),
(status = 500, body = SantiError)
)
)]
pub(super) async fn poll_im(
State(service): State<SantiService>,
Path(participant_id): Path<String>,
Query(params): Query<ImPollParams>,
) -> Result<Json<Vec<ImInboxEntry>>, ApiError> {
service
.im_poll(&participant_id, params.since.unwrap_or(0))
.map(Json)
.map_err(ApiError::from_service)
}

#[derive(serde::Deserialize)]
pub(super) struct ImPollParams {
since: Option<i64>,
}
4 changes: 2 additions & 2 deletions crates/santi-api/src/server/openapi.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,8 +30,8 @@ use utoipa::OpenApi;
super::errors::errors,
super::sse::error_events,
super::routes::runtime_snapshot,
super::routes::send_im,
super::routes::poll_im,
super::im::send_im,
super::im::poll_im,
crate::bucket::get_bucket_object
),
components(schemas(
Expand Down
71 changes: 5 additions & 66 deletions crates/santi-api/src/server/routes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,9 +8,9 @@ use axum::{
use santi_core::{
CompactExecRequest, CompactExecResponse, CompactQueryResponse, CreateSoulRequest,
CreateStrandResponse, CreateWebhookRequest, DriveStrandResponse, ForkStrandResponse,
HealthResponse, ImInboxEntry, ImSendRequest, ImSendResponse, IngestOutcome, MaterialRequest,
SantiError, SantiService, SendStrandAcceptedResponse, SendStrandRequest, Soul, Strand,
StrandBudgetSnapshot, StrandDetail, StrandMaterial, StrandRuntimeSnapshot, WebhookSubscription,
HealthResponse, MaterialRequest, SantiError, SantiService, SendStrandAcceptedResponse,
SendStrandRequest, Soul, Strand, StrandBudgetSnapshot, StrandDetail, StrandMaterial,
StrandRuntimeSnapshot, WebhookSubscription,
};
use tower_http::{
cors::{Any, CorsLayer},
Expand Down Expand Up @@ -50,8 +50,8 @@ pub(super) fn router(service: SantiService) -> Router {
.route("/api/v1/strands/{strand_id}/runtime", get(runtime_snapshot))
// IM layer (orthogonal to the runtime; shares the server for cold-start):
// send into a soul's IM conversation, poll a participant's passive inbox.
.route("/api/v1/im/send", post(send_im))
.route("/api/v1/im/inbox/{participant_id}", get(poll_im))
.route("/api/v1/im/send", post(super::im::send_im))
.route("/api/v1/im/inbox/{participant_id}", get(super::im::poll_im))
.route(
"/api/v1/bucket/{soul_id}/{strand_id}/{*key}",
get(crate::bucket::get_bucket_object),
Expand Down Expand Up @@ -439,67 +439,6 @@ pub(super) async fn strand_budget(
.ok_or_else(|| ApiError::not_found("strand not found"))
}

// ── IM layer routes ─────────────────────────────────────────────────────────
// The plain IM integrated into santi. `strand send`/the runtime stay source-less;
// the participant address is IM envelope only. Inbound reuses the runtime primitive
// (Text into an `im:<participant>` conversation strand). The reply comes back into
// the participant's passive inbox (written by the soul's offline `im reply` egress).

#[utoipa::path(
post,
path = "/api/v1/im/send",
request_body = ImSendRequest,
responses(
(status = 200, body = ImSendResponse),
(status = 423, body = SantiError),
(status = 404, body = SantiError),
(status = 500, body = SantiError)
)
)]
pub(super) async fn send_im(
State(service): State<SantiService>,
Json(request): Json<ImSendRequest>,
) -> Result<Json<ImSendResponse>, ApiError> {
let outcome = service
.im_send(&request.soul_id, &request.participant_id, &request.content)
.map_err(ApiError::from_service)?;
match outcome {
IngestOutcome::Accepted { receipt } => Ok(Json(ImSendResponse {
participant_id: request.participant_id,
receipt,
})),
IngestOutcome::Rejected { error } => Err(ApiError::from_santi(*error)),
}
}

#[utoipa::path(
get,
path = "/api/v1/im/inbox/{participant_id}",
params(
("participant_id" = String, Path),
("since" = Option<i64>, Query)
),
responses(
(status = 200, body = Vec<ImInboxEntry>),
(status = 500, body = SantiError)
)
)]
pub(super) async fn poll_im(
State(service): State<SantiService>,
Path(participant_id): Path<String>,
Query(params): Query<ImPollParams>,
) -> Result<Json<Vec<ImInboxEntry>>, ApiError> {
service
.im_poll(&participant_id, params.since.unwrap_or(0))
.map(Json)
.map_err(ApiError::from_service)
}

#[derive(serde::Deserialize)]
pub(super) struct ImPollParams {
since: Option<i64>,
}

pub(super) async fn openapi() -> Json<utoipa::openapi::OpenApi> {
Json(super::openapi::document())
}
55 changes: 32 additions & 23 deletions crates/santi-api/src/upgrade/finalize.rs
Original file line number Diff line number Diff line change
Expand Up @@ -77,29 +77,7 @@ pub fn finalize_at(

match &request.terminal {
UpgradeTerminal::Upgraded { readiness } => {
store
.resolve_error_incident(
UPGRADE_INCIDENT_KEY,
"upgrade.succeeded",
json!({
"attempt_id": request.attempt_id,
"artifact": bounded_detail(&request.deb),
"terminal": if matches!(readiness, super::UpgradeReadiness::Degraded) {
"upgraded_degraded"
} else {
"upgraded"
},
"readiness": readiness,
}),
)
.map_err(|error| {
Box::new(persistence_error(
&request.attempt_id,
&request.deb,
"upgrade.finalize.resolve_execution",
error,
))
})?;
resolve_upgrade(&store, &request, *readiness)?;
}
UpgradeTerminal::RolledBack { failure } | UpgradeTerminal::Failed { failure } => {
let terminal = if matches!(request.terminal, UpgradeTerminal::RolledBack { .. }) {
Expand All @@ -125,6 +103,37 @@ pub fn finalize_at(
finalize_handover(paths, &store, request, errors)
}

fn resolve_upgrade(
store: &santi_core::SantiStore,
request: &UpgradeFinalizeRequest,
readiness: super::UpgradeReadiness,
) -> Result<(), Box<santi_core::SantiError>> {
store
.resolve_error_incident(
UPGRADE_INCIDENT_KEY,
"upgrade.succeeded",
json!({
"attempt_id": request.attempt_id,
"artifact": bounded_detail(&request.deb),
"terminal": if matches!(readiness, super::UpgradeReadiness::Degraded) {
"upgraded_degraded"
} else {
"upgraded"
},
"readiness": readiness,
}),
)
.map(|_| ())
.map_err(|error| {
Box::new(persistence_error(
&request.attempt_id,
&request.deb,
"upgrade.finalize.resolve_execution",
error,
))
})
}

fn open_execution_failure(
store: &santi_core::SantiStore,
scope: &santi_core::ErrorScope,
Expand Down
2 changes: 1 addition & 1 deletion crates/santi-api/tests/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -144,7 +144,7 @@ async fn send_rejection_locks() {
}

#[tokio::test]
async fn accepted_drive_failure_degrades_health_until_explicit_recovery() {
async fn drive_failure_http_recovery() {
let temp = tempfile::tempdir().expect("temp dir");
let database_path = temp.path().join("santi.sqlite");
let service = SantiService::open(
Expand Down
2 changes: 1 addition & 1 deletion crates/santi-api/tests/upgrade.rs
Original file line number Diff line number Diff line change
Expand Up @@ -185,7 +185,7 @@ fn success_orders_steps() {
}

#[test]
fn degraded_upgrade_does_not_roll_back() {
fn degraded_skips_rollback() {
let mut host = FakeHost {
probe_result: Ok(UpgradeReadiness::Degraded),
..Default::default()
Expand Down
2 changes: 1 addition & 1 deletion crates/santi-core/src/service/flow/ingest.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
use crate::store::{StartTurnOutcome, drive::DriveFailureInput};
use crate::store::{DriveFailureInput, StartTurnOutcome};
use crate::{
DriveStrandResponse, DriveStrandState, ErrorScope, ErrorSource, InboxSource, IngestOutcome,
MessageContent, MessageKind, SantiError, SantiStreamPayload, SendStrandAcceptedResponse,
Expand Down
2 changes: 1 addition & 1 deletion crates/santi-core/src/store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,6 @@ mod assembly;
pub(crate) mod budget;
mod compact;
mod db;
pub(crate) mod drive;
mod errors;
mod fork;
mod im;
Expand All @@ -23,6 +22,7 @@ mod turns;

use db::*;
pub use db::{read_schema_version, soul_memory_file};
pub(crate) use errors::drive::DriveFailureInput;
use rows::{actor_type_db, collect_rows, map_webhook_row, message_state_db};

/// The schema version this binary expects. On open, recognized runtime-schema
Expand Down
4 changes: 2 additions & 2 deletions crates/santi-core/src/store/budget.rs
Original file line number Diff line number Diff line change
Expand Up @@ -93,7 +93,7 @@ impl SantiStore {
.map_err(|error| error.to_string())?;

if let Some(error) =
super::drive::repeat_active_in_conn(&tx, strand_id, "ingest_active_guard")?
super::errors::drive::repeat_active_in_conn(&tx, strand_id, "ingest_active_guard")?
{
tx.commit().map_err(|error| error.to_string())?;
return Ok(IngestOutcome::Rejected {
Expand Down Expand Up @@ -300,7 +300,7 @@ impl SantiStore {
params![turn_id, strand_id, trigger_type, trigger_ref, now],
)
.map_err(|error| error.to_string())?;
super::drive::resolve_in_conn(&tx, strand_id, &turn_id, drained_messages.len())?;
super::errors::drive::resolve_in_conn(&tx, strand_id, &turn_id, drained_messages.len())?;
tx.commit().map_err(|error| error.to_string())?;
Ok(StartTurnOutcome::Started(StartedTurn {
turn: turn_by_id(&conn, &turn_id)?.ok_or_else(|| "created turn missing".to_string())?,
Expand Down
Loading