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
194 changes: 194 additions & 0 deletions crates/arcand/src/canonical.rs
Original file line number Diff line number Diff line change
Expand Up @@ -359,6 +359,48 @@ struct ProviderResponse {
available: Vec<String>,
}

// ─── Consciousness endpoints (BRO-456) ──────────────────────────────────────

/// Request body for POST /sessions/{id}/messages — lightweight message push.
#[derive(Debug, Deserialize, ToSchema)]
struct PushMessageRequest {
/// The message content to send.
message: String,
/// How this message should interact with an active run (default: Collect).
steering: Option<String>,
}

/// Response for POST /sessions/{id}/messages.
#[derive(Debug, Serialize, ToSchema)]
struct PushMessageResponse {
#[schema(value_type = String)]
session_id: SessionId,
queued: bool,
queue_position: Option<usize>,
}

/// Response for GET /sessions/{id}/queue.
#[derive(Debug, Serialize, ToSchema)]
struct QueueStatusResponse {
#[schema(value_type = String)]
session_id: SessionId,
mode: String,
queue_depth: usize,
has_active_run: bool,
oldest_message_age_ms: Option<u64>,
pending: Vec<QueuePendingEntry>,
}

/// A single pending entry in the queue response.
#[derive(Debug, Serialize, ToSchema)]
struct QueuePendingEntry {
id: String,
steering_mode: String,
content: String,
}

// ─── End consciousness types ─────────────────────────────────────────────────

#[derive(Debug, Serialize, ToSchema)]
struct RunResponse {
#[schema(value_type = String, example = "sess-abc123")]
Expand Down Expand Up @@ -452,6 +494,8 @@ struct ErrorResponse {
get_context,
get_cost,
list_mcp_servers,
push_message,
get_queue_status,
),
components(schemas(
// Request / response types
Expand All @@ -477,6 +521,11 @@ struct ErrorResponse {
AutonomicResponse,
ContextResponse,
CostResponse,
// Consciousness endpoint types (BRO-456)
PushMessageRequest,
PushMessageResponse,
QueueStatusResponse,
QueuePendingEntry,
// Mirror schemas for external aios-protocol types
OperatingModeSchema,
AgentStateVectorSchema,
Expand All @@ -499,6 +548,7 @@ struct ErrorResponse {
(name = "provider", description = "Live provider switching"),
(name = "autonomic", description = "Autonomic context regulation and homeostatic state"),
(name = "mcp", description = "MCP server registry"),
(name = "consciousness", description = "Consciousness session management and message queue"),
)
)]
struct ApiDoc;
Expand Down Expand Up @@ -740,6 +790,15 @@ pub fn create_canonical_router_with_skills(
"/sessions/{session_id}/approvals/{approval_id}",
post(resolve_approval),
)
// BRO-456: Consciousness message push and queue introspection.
.route(
"/sessions/{session_id}/messages",
post(push_message),
)
.route(
"/sessions/{session_id}/queue",
get(get_queue_status),
)
.route("/provider", get(get_provider).put(set_provider))
.route("/autonomic", get(get_autonomic))
.route("/context", get(get_context))
Expand Down Expand Up @@ -2641,6 +2700,141 @@ async fn get_cost(State(state): State<CanonicalState>) -> Json<CostResponse> {
})
}

// ─── Consciousness endpoints (BRO-456) ──────────────────────────────────────

/// Push a message to an active consciousness session.
///
/// Lightweight alternative to `POST /runs` for follow-ups, interrupts, or steers.
/// Requires consciousness mode (`ARCAN_CONSCIOUSNESS=true`).
#[utoipa::path(
post,
path = "/sessions/{session_id}/messages",
tag = "consciousness",
params(("session_id" = String, Path, description = "Session identifier")),
request_body = PushMessageRequest,
responses(
(status = 202, description = "Message accepted", body = PushMessageResponse),
(status = 404, description = "Session not found or consciousness disabled", body = ErrorResponse),
(status = 503, description = "Message rejected", body = ErrorResponse)
)
)]
async fn push_message(
Path(session_id): Path<String>,
State(state): State<CanonicalState>,
Json(request): Json<PushMessageRequest>,
) -> Result<(StatusCode, Json<PushMessageResponse>), (StatusCode, Json<serde_json::Value>)> {
let registry = state.consciousness_registry.as_ref().ok_or_else(|| {
(
StatusCode::NOT_FOUND,
Json(json!({ "error": "consciousness mode is not enabled" })),
)
})?;

let handle = registry.get(&session_id).ok_or_else(|| {
(
StatusCode::NOT_FOUND,
Json(json!({ "error": "no active consciousness session", "session_id": session_id })),
)
})?;

let steering = match request.steering.as_deref() {
Some("steer" | "Steer") => aios_protocol::SteeringMode::Steer,
Some("interrupt" | "Interrupt") => aios_protocol::SteeringMode::Interrupt,
Some("followup" | "Followup") => aios_protocol::SteeringMode::Followup,
_ => aios_protocol::SteeringMode::Collect,
};

let (ack_tx, ack_rx) = tokio::sync::oneshot::channel();
let send_result = handle
.send(crate::consciousness::ConsciousnessEvent::UserMessage(
Box::new(crate::consciousness::UserMessageEvent {
objective: request.message,
branch: BranchId::main(),
steering,
ack: Some(ack_tx),
run_context: crate::consciousness::RunContext::default(),
}),
))
.await;

if send_result.is_err() {
return Err(internal_error(format!(
"consciousness actor for session {} is not running",
session_id
)));
}

let sid = SessionId::from_string(session_id);
match tokio::time::timeout(Duration::from_secs(5), ack_rx).await {
Ok(Ok(crate::consciousness::ConsciousnessAck::Accepted { queued })) => Ok((
StatusCode::ACCEPTED,
Json(PushMessageResponse {
session_id: sid,
queued,
queue_position: if queued { Some(1) } else { None },
}),
)),
Comment on lines +2769 to +2776

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟡 Minor

Hardcoded queue_position: Some(1) is inaccurate.

When queued is true, the response always reports queue_position: Some(1), but this doesn't reflect the actual position in the queue. The ConsciousnessAck::Accepted variant doesn't provide the queue position. Either remove this field, make it None when queued, or extend the ConsciousnessAck to include the actual position.

🐛 Option: Return None for queue_position until accurate data is available
         Ok(Ok(crate::consciousness::ConsciousnessAck::Accepted { queued })) => Ok((
             StatusCode::ACCEPTED,
             Json(PushMessageResponse {
                 session_id: sid,
                 queued,
-                queue_position: if queued { Some(1) } else { None },
+                queue_position: None, // TODO: Extend ConsciousnessAck to include actual position
             }),
         )),
📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
Ok(Ok(crate::consciousness::ConsciousnessAck::Accepted { queued })) => Ok((
StatusCode::ACCEPTED,
Json(PushMessageResponse {
session_id: sid,
queued,
queue_position: if queued { Some(1) } else { None },
}),
)),
Ok(Ok(crate::consciousness::ConsciousnessAck::Accepted { queued })) => Ok((
StatusCode::ACCEPTED,
Json(PushMessageResponse {
session_id: sid,
queued,
queue_position: None, // TODO: Extend ConsciousnessAck to include actual position
}),
)),
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@crates/arcand/src/canonical.rs` around lines 2769 - 2776, The response
currently hardcodes queue_position: Some(1) for ConsciousnessAck::Accepted which
is inaccurate; update the match arm handling
crate::consciousness::ConsciousnessAck::Accepted in canonical.rs so that when
queued is true the PushMessageResponse sets queue_position to None (or simply
always set queue_position: None) instead of Some(1); if you intend to return a
real position later, extend the ConsciousnessAck::Accepted variant to carry the
position and map that value into PushMessageResponse.queue_position instead.

Ok(Ok(crate::consciousness::ConsciousnessAck::Rejected { reason })) => Err((
StatusCode::SERVICE_UNAVAILABLE,
Json(json!({ "error": "rejected", "reason": reason })),
)),
_ => Err(internal_error("consciousness actor did not respond")),
}
}

/// Get the queue status for a consciousness session.
#[utoipa::path(
get,
path = "/sessions/{session_id}/queue",
tag = "consciousness",
params(("session_id" = String, Path, description = "Session identifier")),
responses(
(status = 200, description = "Queue status", body = QueueStatusResponse),
(status = 404, description = "Session not found or consciousness disabled", body = ErrorResponse)
)
)]
async fn get_queue_status(
Path(session_id): Path<String>,
State(state): State<CanonicalState>,
) -> Result<Json<QueueStatusResponse>, (StatusCode, Json<serde_json::Value>)> {
let registry = state.consciousness_registry.as_ref().ok_or_else(|| {
(
StatusCode::NOT_FOUND,
Json(json!({ "error": "consciousness mode is not enabled" })),
)
})?;

let handle = registry.get(&session_id).ok_or_else(|| {
(
StatusCode::NOT_FOUND,
Json(json!({ "error": "no active consciousness session", "session_id": session_id })),
)
})?;

let status = handle
.query_status()
.await
.ok_or_else(|| internal_error("consciousness actor did not respond to status query"))?;

let sid = SessionId::from_string(session_id);
Ok(Json(QueueStatusResponse {
session_id: sid,
mode: status.mode,
queue_depth: status.queue_depth,
has_active_run: status.has_active_run,
oldest_message_age_ms: status.oldest_message_age_ms,
pending: status
.queue_pending
.into_iter()
.map(|m| QueuePendingEntry {
id: m.id,
steering_mode: m.mode,
content: m.content,
})
.collect(),
}))
}

// ─── Error helpers ───────────────────────────────────────────────────────────

fn internal_error(error: impl std::fmt::Display) -> (StatusCode, Json<serde_json::Value>) {
Expand Down
63 changes: 63 additions & 0 deletions crates/arcand/src/consciousness.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,12 +53,32 @@ impl Default for ConsciousnessConfig {
pub enum ConsciousnessEvent {
/// User message from HTTP POST /runs or /messages.
UserMessage(Box<UserMessageEvent>),
/// Query the actor's current status (mode + queue snapshot).
QueryStatus { reply: oneshot::Sender<ActorStatus> },
/// Timer tick from internal intervals.
TimerTick { tick_type: TimerTickType },
/// Graceful shutdown request.
Shutdown,
}

/// Snapshot of the actor's status for the GET /queue endpoint.
#[derive(Debug, Clone, serde::Serialize)]
pub struct ActorStatus {
pub mode: String,
pub queue_depth: usize,
pub queue_pending: Vec<PendingMessage>,
pub has_active_run: bool,
pub oldest_message_age_ms: Option<u64>,
}

/// A pending message in the queue (serializable for API responses).
#[derive(Debug, Clone, serde::Serialize)]
pub struct PendingMessage {
pub id: String,
pub mode: String,
pub content: String,
}

/// Payload for a user message event (boxed to keep enum small).
#[derive(Debug)]
pub struct UserMessageEvent {
Expand Down Expand Up @@ -211,6 +231,10 @@ impl SessionConsciousness {
)
.await;
}
ConsciousnessEvent::QueryStatus { reply } => {
let status = self.build_status();
let _ = reply.send(status);
}
ConsciousnessEvent::TimerTick { tick_type } => {
self.handle_timer_tick(tick_type).await;
}
Expand Down Expand Up @@ -497,6 +521,32 @@ impl SessionConsciousness {
}
}
}

/// Build a status snapshot for the QueryStatus response.
fn build_status(&self) -> ActorStatus {
let queue_status = self.state.queue.status().ok();
let pending = queue_status
.as_ref()
.map(|s| {
s.pending
.iter()
.map(|m| PendingMessage {
id: m.id.clone(),
mode: format!("{:?}", m.mode),
content: m.content.clone(),
})
.collect()
})
.unwrap_or_default();

ActorStatus {
mode: format!("{:?}", self.state.mode),
queue_depth: queue_status.as_ref().map(|s| s.depth).unwrap_or(0),
queue_pending: pending,
has_active_run: queue_status.as_ref().is_some_and(|s| s.has_active_run),
oldest_message_age_ms: queue_status.and_then(|s| s.oldest_message_age_ms),
}
}
}

// ─── Handle ─────────────────────────────────────────────────────────────────
Expand Down Expand Up @@ -534,6 +584,19 @@ impl ConsciousnessHandle {
self.tx.send(event).await
}

/// Query the actor's current status (mode + queue snapshot).
pub async fn query_status(&self) -> Option<ActorStatus> {
let (reply_tx, reply_rx) = oneshot::channel();
self.tx
.send(ConsciousnessEvent::QueryStatus { reply: reply_tx })
.await
.ok()?;
tokio::time::timeout(Duration::from_secs(2), reply_rx)
.await
.ok()?
.ok()
}

/// Send a shutdown event and wait for the actor to stop.
pub async fn shutdown(self) {
let _ = self.tx.send(ConsciousnessEvent::Shutdown).await;
Expand Down
Loading
Loading