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: 1 addition & 0 deletions crates/forge-schemas/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,5 +10,6 @@ pub mod handoff;
pub mod ledger;
pub mod telemetry;
pub mod verification;
pub mod workflow;

pub use canonical::{canonical_json, request_hash, sha256_hex};
127 changes: 127 additions & 0 deletions crates/forge-schemas/src/workflow.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,127 @@
//! Grafo do Squad Designer (`squad.workflow.v1`, Fase 7 Onda 14).
//!
//! Salvar valida a forma (schema + integridade de arestas) e grava no
//! ledger — **não aplica** ao orquestrador real: o `UnifiedOrchestrator`
//! continua com os 5 agentes fixos (`forge_squad`), sem reescrita nesta
//! fase. "Salvar honesto": o servidor confirma que o grafo foi validado e
//! persistido, nunca que o squad passou a usá-lo.

use schemars::JsonSchema;
use serde::{Deserialize, Serialize};

#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum WorkflowNodeKind {
Card,
Pill,
}

#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
pub struct WorkflowNodeParam {
pub k: String,
pub v: String,
}

#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
pub struct WorkflowNode {
pub id: String,
pub x: f64,
pub y: f64,
pub kind: WorkflowNodeKind,
pub name: String,
pub role: String,
pub color: String,
pub icon: String,
pub sub: String,
pub params: Vec<WorkflowNodeParam>,
pub removable: bool,
}

#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
pub struct WorkflowEdge {
pub from: String,
pub to: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub label: Option<String>,
}

#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
pub struct SquadWorkflow {
pub nodes: Vec<WorkflowNode>,
pub edges: Vec<WorkflowEdge>,
}

impl SquadWorkflow {
/// Única checagem semântica além do schema (campos/tipos obrigatórios já
/// cobertos por serde + JSON Schema): toda aresta referencia um nó que
/// existe. Erro aponta o lado (`from`/`to`) e o id que falhou — um 422
/// claro, não um 500 genérico nem um grafo salvo silenciosamente
/// quebrado.
pub fn validate_edges(&self) -> Result<(), String> {
let ids: std::collections::HashSet<&str> =
self.nodes.iter().map(|n| n.id.as_str()).collect();
for edge in &self.edges {
if !ids.contains(edge.from.as_str()) {
return Err(format!(
"aresta referencia nó inexistente em 'from': {}",
edge.from
));
}
if !ids.contains(edge.to.as_str()) {
return Err(format!(
"aresta referencia nó inexistente em 'to': {}",
edge.to
));
}
}
Ok(())
}
}

#[cfg(test)]
mod tests {
use super::*;

fn node(id: &str) -> WorkflowNode {
WorkflowNode {
id: id.into(),
x: 0.0,
y: 0.0,
kind: WorkflowNodeKind::Card,
name: id.into(),
role: "agente".into(),
color: "var(--rust)".into(),
icon: "◆".into(),
sub: "".into(),
params: vec![],
removable: true,
}
}

#[test]
fn grafo_com_arestas_validas_passa() {
let wf = SquadWorkflow {
nodes: vec![node("a"), node("b")],
edges: vec![WorkflowEdge {
from: "a".into(),
to: "b".into(),
label: None,
}],
};
assert!(wf.validate_edges().is_ok());
}

#[test]
fn aresta_para_no_inexistente_e_rejeitada_com_erro_claro() {
let wf = SquadWorkflow {
nodes: vec![node("a")],
edges: vec![WorkflowEdge {
from: "a".into(),
to: "fantasma".into(),
label: None,
}],
};
let err = wf.validate_edges().unwrap_err();
assert!(err.contains("fantasma"), "erro deveria citar o id: {err}");
}
}
27 changes: 27 additions & 0 deletions crates/forge-schemas/tests/schema_fixtures.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ use forge_schemas::experiment::ExperimentReport;
use forge_schemas::handoff::HandoffEvent;
use forge_schemas::ledger::LedgerEntry;
use forge_schemas::telemetry::TelemetryEvent;
use forge_schemas::workflow::SquadWorkflow;
use jsonschema::validator_for;
use serde_json::Value;

Expand Down Expand Up @@ -128,6 +129,32 @@ fn experiment_fixture_valida_e_desserializa() {
);
}

/// A checagem semântica (aresta referencia nó inexistente) não é
/// expressável em JSON Schema puro — fica em `SquadWorkflow::validate_edges`
/// (testada isoladamente em `workflow.rs`). Aqui só a FORMA: campo
/// obrigatório ausente (`removable`) deve reprovar o schema.
#[test]
fn squad_workflow_fixture_valida_e_desserializa() {
let schema = schema("squad-workflow");
let doc = fixture("squad-workflow");
let validator = validator_for(&schema).expect("schema compila");

assert!(
validator.is_valid(&doc["valid"]),
"fixture válida não bateu o schema: {:?}",
validator.iter_errors(&doc["valid"]).collect::<Vec<_>>()
);
let parsed: SquadWorkflow =
serde_json::from_value(doc["valid"].clone()).expect("desserializa em SquadWorkflow");
assert_eq!(parsed.nodes.len(), 2);
assert!(parsed.validate_edges().is_ok());

assert!(
!validator.is_valid(&doc["invalid_missing_removable"]),
"documento sem 'removable' deveria reprovar o schema"
);
}

/// Sem tipo Rust/Python — só protege o schema em si (sintaxe/drift), não
/// uma paridade de tipo. Ver nota no topo do arquivo e no `$comment` da
/// fixture.
Expand Down
162 changes: 162 additions & 0 deletions crates/forge-server/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ use axum::Router;
use forge_llm::model_tier::{tier_from_id, ModelTier};
use forge_llm::rate_limit::RateLimiter;
use forge_schemas::experiment::{ExperimentReport, VariantStats};
use forge_schemas::workflow::SquadWorkflow;
use forge_store::{LedgerStore, PromptLibrary, Telemetry};
use serde::{Deserialize, Serialize};
use serde_json::Value;
Expand Down Expand Up @@ -119,6 +120,7 @@ pub fn router(
.route("/api/providers", get(list_providers))
.route("/api/verify/run", post(run_verify_start))
.route("/api/verify/{id}", get(get_verify_status))
.route("/api/designer/workflow", post(save_workflow))
.fallback_service(serve_dir)
.with_state(AppState {
telemetry,
Expand Down Expand Up @@ -656,6 +658,60 @@ async fn verify_ledger(State(state): State<AppState>) -> Response {
}
}

#[derive(Serialize)]
struct SaveWorkflowResponse {
seq: u64,
workflow_id: &'static str,
}

/// `POST /api/designer/workflow` (Fase 7 Onda 14) — valida o grafo do Squad
/// Designer contra `squad.workflow.v1` (schema + integridade de arestas via
/// `SquadWorkflow::validate_edges`) e grava no ledger (mesmo
/// `LedgerStore::append` que toda outra escrita de auditoria já usa — zero
/// mudança de ledger). "Salvar honesto": confirma que o grafo foi validado e
/// persistido, nunca que o orquestrador passou a usá-lo — os 5 agentes
/// fixos do `UnifiedOrchestrator` continuam decidindo, sem reescrita nesta
/// fase (aplicar o grafo real é trabalho futuro).
async fn save_workflow(
State(state): State<AppState>,
Json(workflow): Json<SquadWorkflow>,
) -> Response {
if let Err(e) = workflow.validate_edges() {
return (
StatusCode::UNPROCESSABLE_ENTITY,
Json(ErrorBody::new("invalid_workflow", e)),
)
.into_response();
}
let payload = match serde_json::to_value(&workflow) {
Ok(v) => v,
Err(e) => return db_error(e),
};
let entry = forge_schemas::ledger::LedgerEntry {
seq: 0,
prev_hash: String::new(),
entry_hash: String::new(),
kind: "designer.workflow_saved".into(),
actor: "web:designer".into(),
payload,
r#override: None,
fake_marker: None,
ts: now_rfc3339(),
};
let mut ledger = state.ledger.lock().unwrap_or_else(|e| e.into_inner());
match ledger.append(entry) {
Ok(saved) => (
StatusCode::CREATED,
Json(SaveWorkflowResponse {
seq: saved.seq,
workflow_id: "squad.workflow.v1",
}),
)
.into_response(),
Err(e) => db_error(e),
}
}

#[cfg(test)]
mod tests {
use super::*;
Expand Down Expand Up @@ -1262,6 +1318,112 @@ mod tests {
assert!(json.get("error").is_none());
}

fn workflow_body(edges: serde_json::Value) -> serde_json::Value {
serde_json::json!({
"nodes": [
{
"id": "task", "x": 0.0, "y": 0.0, "kind": "pill", "name": "tarefa",
"role": "entrada", "color": "c", "icon": "▸", "sub": "", "params": [],
"removable": false
},
{
"id": "architect", "x": 10.0, "y": 10.0, "kind": "card", "name": "architect",
"role": "arquitetura", "color": "c", "icon": "◆", "sub": "", "params": [],
"removable": true
}
],
"edges": edges,
})
}

/// Fronteira da Onda 14 (Designer, "salvar honesto"): grafo válido grava
/// no MESMO ledger que a rota de leitura já usa — lido direto de volta
/// (não uma segunda fonte de verdade), `seq` real (não fabricado no
/// cliente), `kind`/`actor` corretos.
#[tokio::test]
async fn salvar_workflow_valido_grava_no_ledger_e_e_lido_de_volta() {
let ledger = ledger_vazio();
let web_dir = fixture_web_dir();
let app = router(
Telemetry::open_in_memory().unwrap(),
prompt_library_vazia(),
Arc::clone(&ledger),
web_dir.path(),
web_dir.path(),
);

let resp = app
.oneshot(
Request::builder()
.method("POST")
.uri("/api/designer/workflow")
.header("content-type", "application/json")
.body(Body::from(
workflow_body(serde_json::json!([{"from": "task", "to": "architect"}]))
.to_string(),
))
.unwrap(),
)
.await
.unwrap();
assert_eq!(resp.status(), StatusCode::CREATED);
let body = axum::body::to_bytes(resp.into_body(), usize::MAX)
.await
.unwrap();
let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
assert_eq!(json["seq"], 1);
assert_eq!(json["workflow_id"], "squad.workflow.v1");

// Lido direto do MESMO storage por trás da rota — não uma segunda
// cópia inventada na resposta HTTP.
let store = ledger.lock().unwrap();
let entries = store.recent(10, None).unwrap();
assert_eq!(entries.len(), 1);
assert_eq!(entries[0].kind, "designer.workflow_saved");
assert_eq!(entries[0].actor, "web:designer");
assert_eq!(entries[0].payload["nodes"].as_array().unwrap().len(), 2);
}

/// Grafo malformado (aresta pra nó inexistente) é rejeitado com erro
/// claro (422, citando o id) — não salvo silenciosamente. O ledger
/// continua vazio: a validação acontece ANTES do `append`, não depois.
#[tokio::test]
async fn salvar_workflow_com_aresta_pendente_e_rejeitado_e_nao_grava_nada() {
let ledger = ledger_vazio();
let web_dir = fixture_web_dir();
let app = router(
Telemetry::open_in_memory().unwrap(),
prompt_library_vazia(),
Arc::clone(&ledger),
web_dir.path(),
web_dir.path(),
);

let resp = app
.oneshot(
Request::builder()
.method("POST")
.uri("/api/designer/workflow")
.header("content-type", "application/json")
.body(Body::from(
workflow_body(serde_json::json!([{"from": "task", "to": "fantasma"}]))
.to_string(),
))
.unwrap(),
)
.await
.unwrap();
assert_eq!(resp.status(), StatusCode::UNPROCESSABLE_ENTITY);
let body = axum::body::to_bytes(resp.into_body(), usize::MAX)
.await
.unwrap();
let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
assert!(json["error"].as_str().unwrap().contains("fantasma"));

let store = ledger.lock().unwrap();
assert_eq!(store.recent(10, None).unwrap().len(), 0);
}

/// Fronteira da Onda 7 (A5): `GET /api/models/usage` bate por igualdade
/// com agregação MANUAL dos mesmos eventos semeados — inclui a coluna
/// `tier` derivada de `tier_from_id` (não fabricada), e não conta um
Expand Down
Loading
Loading