feat(umind): vigilancia de APIs — avisa cuando pasa algo, sin quemar plata
"Avisame cuando esta API devuelva stock en cero." La vigilancia apunta a una HERRAMIENTA ya configurada del agente y nunca a una URL libre: así hereda entero el blindaje de LlamarHerramientaWebhook —DNS resuelto en el connect, IPs internas rechazadas, sin redirects, respuesta acotada— sin reimplementar una línea de eso. Hay un test que falla si alguien le agrega un campo url. El ahorro que hace viable la feature: se hashea la respuesta y, si no cambió, no hay nada que evaluar y no se gasta un token. Preguntarle a la IA en cada chequeo serían 1440 consultas diarias por cliente para responder algo que ya sabíamos. El test verifica que el corte por hash esté ANTES de la evaluación. Y se avisa al ENTRAR en condición, no en cada chequeo que la siga cumpliendo: una alerta que llega cada media hora deja de leerse a la segunda. Es la lección que el monitor de sitios ya había aprendido con las transiciones de estado. Solo un SI explícito dispara el aviso. Un modelo que devuelve una explicación, un error o cualquier otra cosa se toma como que no: un aviso que no llega molesta menos que uno que despierta a alguien a las 3 AM sin motivo. Si el envío del aviso falla, no se marca como avisado — se reintenta al siguiente ciclo. Y tras diez fallos seguidos la vigilancia se apaga sola: una API que dejó de existir no puede consultarse para siempre ni seguir gastando el plan de quien la configuró. Máximo 5 por agente, mínimo 5 minutos de intervalo. Solo canal interno: cada vigilancia es una llamada saliente recurrente y consultas a la IA, o sea plan del dueño gastado por quien la pida. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
047837dd23
commit
7ddc3f5286
@@ -85,6 +85,12 @@ func IniciarCron() {
|
||||
return
|
||||
}
|
||||
|
||||
// Vigilancias de APIs — cada una según su propio intervalo.
|
||||
if _, err := cronScheduler.AddFunc("* * * * *", RevisarVigilancias); err != nil {
|
||||
log.Printf("[CRON] Error registrando tarea vigilancias_umind: %v", err)
|
||||
return
|
||||
}
|
||||
|
||||
// Los pendientes que nadie miró vencen solos.
|
||||
if _, err := cronScheduler.AddFunc("0 5 * * *", func() { models.VencerAccionesViejas() }); err != nil {
|
||||
log.Printf("[CRON] Error registrando tarea vencer_acciones: %v", err)
|
||||
|
||||
@@ -68,6 +68,7 @@ func umindTools(agenteID uint, sesionInterna bool) []agentTool {
|
||||
// cualquiera, y programar recordatorios gasta el plan de otro.
|
||||
if sesionInterna {
|
||||
tools = append(tools, umindAvisoTools()...)
|
||||
tools = append(tools, umindVigilanciaTools(agenteID)...)
|
||||
}
|
||||
|
||||
herramientas, err := models.GetUmindHerramientasActivas(agenteID)
|
||||
@@ -190,6 +191,13 @@ func ejecutarHerramienta(agenteID uint, sessionID string, sesionInterna bool, na
|
||||
return executeUmindDocumentoTool(agenteID, args)
|
||||
}
|
||||
|
||||
if strings.HasSuffix(name, "_vigilancia") || name == "listar_vigilancias" {
|
||||
if !sesionInterna {
|
||||
return `{"error": "las vigilancias solo puede crearlas el dueño desde su canal privado"}`
|
||||
}
|
||||
return executeUmindVigilanciaTool(agenteID, sessionID, name, args)
|
||||
}
|
||||
|
||||
if strings.HasSuffix(name, "_aviso") || name == "listar_avisos" {
|
||||
if !sesionInterna {
|
||||
return `{"error": "los recordatorios solo puede programarlos el dueño desde su canal privado"}`
|
||||
|
||||
@@ -0,0 +1,161 @@
|
||||
package services
|
||||
|
||||
import (
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"fmt"
|
||||
"log"
|
||||
"strings"
|
||||
"sync"
|
||||
|
||||
"github.com/sujit-baniya/fiber-boilerplate/pkg/models"
|
||||
)
|
||||
|
||||
// La vigilancia de APIs: "avisame cuando esta API devuelva stock en cero".
|
||||
|
||||
var vigilanciasEnCurso sync.Mutex
|
||||
|
||||
// vigilanciaFallosMax es cuántos errores seguidos aguanta antes de apagarse.
|
||||
// Una API que dejó de existir no puede consultarse para siempre.
|
||||
const vigilanciaFallosMax = 10
|
||||
|
||||
// RevisarVigilancias corre cada minuto y atiende solo las que ya les toca.
|
||||
func RevisarVigilancias() {
|
||||
if !vigilanciasEnCurso.TryLock() {
|
||||
return
|
||||
}
|
||||
defer vigilanciasEnCurso.Unlock()
|
||||
|
||||
vigilancias, err := models.GetUmindVigilanciasPendientes()
|
||||
if err != nil {
|
||||
log.Printf("[UMIND_VIGILA] no se pudieron listar: %v", err)
|
||||
return
|
||||
}
|
||||
if len(vigilancias) == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
const trabajadores = 4
|
||||
cola := make(chan models.UmindVigilancia)
|
||||
var wg sync.WaitGroup
|
||||
for w := 0; w < trabajadores; w++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
for v := range cola {
|
||||
revisarUnaVigilancia(v)
|
||||
}
|
||||
}()
|
||||
}
|
||||
for _, v := range vigilancias {
|
||||
cola <- v
|
||||
}
|
||||
close(cola)
|
||||
wg.Wait()
|
||||
}
|
||||
|
||||
func revisarUnaVigilancia(v models.UmindVigilancia) {
|
||||
herramienta, err := models.GetUmindHerramientaByNombre(v.AgenteID, v.Herramienta)
|
||||
if err != nil {
|
||||
models.DesactivarVigilancia(v.ID, "la herramienta que vigilaba ya no existe")
|
||||
models.RegistrarEventoUmind(v.AgenteID, "warn", "vigilancia",
|
||||
"Se apagó la vigilancia "+v.Nombre+": la herramienta que usaba ya no existe", "")
|
||||
return
|
||||
}
|
||||
authValor := ""
|
||||
if herramienta.AuthHeaderValorEnc != "" {
|
||||
if authValor, err = DescifrarSecretoUmind(herramienta.AuthHeaderValorEnc); err != nil {
|
||||
models.ActualizarVigilancia(v.ID, map[string]interface{}{"ultimo_error": "credencial ilegible"})
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
respuesta, err := LlamarHerramientaWebhook(herramienta.URL, herramienta.AuthHeaderNombre, authValor, map[string]interface{}{})
|
||||
if err != nil {
|
||||
fallos := v.FallosNoOk + 1
|
||||
if fallos >= vigilanciaFallosMax {
|
||||
models.DesactivarVigilancia(v.ID, err.Error())
|
||||
models.RegistrarEventoUmind(v.AgenteID, "warn", "vigilancia",
|
||||
"Se apagó la vigilancia "+v.Nombre+" tras fallar muchas veces seguidas", err.Error())
|
||||
return
|
||||
}
|
||||
models.ActualizarVigilancia(v.ID, map[string]interface{}{"ultimo_error": err.Error(), "fallos_no_ok": fallos})
|
||||
return
|
||||
}
|
||||
|
||||
// Acá está el ahorro: si la respuesta es idéntica a la anterior, no hay
|
||||
// nada que evaluar. Preguntarle a la IA en cada chequeo costaría una
|
||||
// fortuna por una respuesta que ya sabemos.
|
||||
suma := sha256.Sum256([]byte(respuesta))
|
||||
hash := hex.EncodeToString(suma[:])
|
||||
if hash == v.HashUltimo {
|
||||
models.ActualizarVigilancia(v.ID, map[string]interface{}{"ultimo_error": "", "fallos_no_ok": 0})
|
||||
return
|
||||
}
|
||||
|
||||
cumple, err := EvaluarCondicionVigilancia(v.Condicion, respuesta)
|
||||
if err != nil {
|
||||
models.ActualizarVigilancia(v.ID, map[string]interface{}{
|
||||
"hash_ultimo": hash, "ultimo_error": "no se pudo evaluar la condición: " + err.Error(),
|
||||
})
|
||||
return
|
||||
}
|
||||
models.RegistrarUsoUmind(v.AgenteID, models.UsoTipoIA, 400, "tokens")
|
||||
|
||||
updates := map[string]interface{}{"hash_ultimo": hash, "ultimo_error": "", "fallos_no_ok": 0, "en_condicion": cumple}
|
||||
// Se avisa al ENTRAR en condición. Si ya estaba y sigue, no se repite:
|
||||
// una alerta que llega cada media hora deja de leerse a la segunda.
|
||||
if cumple && !v.EnCondicion {
|
||||
if err := avisarVigilancia(v, respuesta); err != nil {
|
||||
log.Printf("[UMIND_VIGILA] no se pudo avisar de %q: %v", v.Nombre, err)
|
||||
models.RegistrarEventoUmind(v.AgenteID, "error", "vigilancia", "No se pudo avisar de la vigilancia "+v.Nombre, err.Error())
|
||||
// Sin marcar en_condicion: se reintenta el aviso al próximo ciclo.
|
||||
delete(updates, "en_condicion")
|
||||
}
|
||||
}
|
||||
models.ActualizarVigilancia(v.ID, updates)
|
||||
}
|
||||
|
||||
// EvaluarCondicionVigilancia le pregunta a la IA si lo que devolvió la API
|
||||
// cumple la condición. Se le exige SI o NO para que la respuesta sea usable:
|
||||
// un modelo que explica su razonamiento acá no sirve de nada.
|
||||
func EvaluarCondicionVigilancia(condicion, respuesta string) (bool, error) {
|
||||
if len(respuesta) > 6000 {
|
||||
respuesta = respuesta[:6000]
|
||||
}
|
||||
sistema := "Sos un evaluador. Te dan la respuesta de una API y una condición. " +
|
||||
"Respondé EXACTAMENTE una palabra: SI si la condición se cumple, NO si no se cumple. " +
|
||||
"Ante la duda respondé NO."
|
||||
usuario := fmt.Sprintf("Condición: %s\n\nRespuesta de la API:\n%s", condicion, respuesta)
|
||||
out, err := CompletarTextoIA("ia", sistema, usuario)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
return interpretarVeredicto(out), nil
|
||||
}
|
||||
|
||||
// interpretarVeredicto acepta solo un SI explícito. Cualquier otra cosa —una
|
||||
// explicación, un error, un modelo confundido— se toma como que no: es
|
||||
// preferible un aviso que no llega a uno que llega sin motivo a las 3 AM.
|
||||
func interpretarVeredicto(out string) bool {
|
||||
limpio := strings.ToUpper(strings.TrimSpace(out))
|
||||
limpio = strings.Trim(limpio, ".!\"' \n")
|
||||
return limpio == "SI" || limpio == "SÍ" ||
|
||||
strings.HasPrefix(limpio, "SI ") || strings.HasPrefix(limpio, "SI,") ||
|
||||
strings.HasPrefix(limpio, "SÍ ") || strings.HasPrefix(limpio, "SÍ,")
|
||||
}
|
||||
|
||||
// avisarVigilancia reusa la entrega de los recordatorios: mismo destino, mismo
|
||||
// camino, una sola implementación que mantener.
|
||||
func avisarVigilancia(v models.UmindVigilancia, respuesta string) error {
|
||||
detalle := v.Condicion
|
||||
if len(respuesta) < 800 {
|
||||
detalle += "\n\nLo que devolvió la API:\n" + respuesta
|
||||
}
|
||||
return entregarAviso(models.UmindAviso{
|
||||
AgenteID: v.AgenteID,
|
||||
Titulo: "Se cumplió lo que estabas vigilando: " + v.Nombre,
|
||||
Detalle: detalle,
|
||||
Destino: v.Destino,
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,72 @@
|
||||
package services
|
||||
|
||||
import (
|
||||
"os"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// Solo un SI explícito dispara el aviso. Un modelo que devuelve una
|
||||
// explicación, un error o cualquier otra cosa no puede despertar a nadie a las
|
||||
// 3 de la mañana: un aviso que no llega molesta menos que uno sin motivo.
|
||||
func TestEvaluarCondicionSoloAceptaSiExplicito(t *testing.T) {
|
||||
casos := map[string]bool{
|
||||
"SI": true, "si": true, "SÍ": true, "Sí.": true,
|
||||
"SI, el stock está en cero": true,
|
||||
"NO": false, "no": false,
|
||||
"Sin duda no puedo determinarlo": false,
|
||||
"La condición parece cumplirse en algunos casos": false,
|
||||
"": false, "ERROR": false, "true": false,
|
||||
}
|
||||
for entrada, quiero := range casos {
|
||||
if got := interpretarVeredicto(entrada); got != quiero {
|
||||
t.Errorf("veredicto %q = %v, quiero %v", entrada, got, quiero)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// El ahorro que hace viable la feature: sin cambio en la respuesta, no se
|
||||
// llama a la IA. Sin esto, cinco vigilancias cada 5 minutos son 1440 consultas
|
||||
// diarias al modelo por cliente.
|
||||
func TestVigilanciaNoLlamaALaIASiNoCambio(t *testing.T) {
|
||||
b, err := os.ReadFile("umind_vigilancia_service.go")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
s := string(b)
|
||||
i := strings.Index(s, "func revisarUnaVigilancia")
|
||||
cuerpo := s[i:]
|
||||
cuerpo = cuerpo[:strings.Index(cuerpo, "\n// EvaluarCondicion")]
|
||||
|
||||
hash := strings.Index(cuerpo, "hash == v.HashUltimo")
|
||||
eval := strings.Index(cuerpo, "EvaluarCondicionVigilancia")
|
||||
if hash < 0 || eval < 0 || hash > eval {
|
||||
t.Error("el corte por hash tiene que ir ANTES de evaluar con la IA")
|
||||
}
|
||||
corte := cuerpo[hash:eval]
|
||||
if !strings.Contains(corte, "return") {
|
||||
t.Error("si el hash no cambió tiene que cortar, no seguir a la evaluación")
|
||||
}
|
||||
|
||||
// Se avisa al entrar en condición, no en cada chequeo que la siga
|
||||
// cumpliendo: una alerta cada media hora deja de leerse a la segunda.
|
||||
if !strings.Contains(cuerpo, "cumple && !v.EnCondicion") {
|
||||
t.Error("falta el flanco: sin él la misma alerta se repite en cada ciclo")
|
||||
}
|
||||
}
|
||||
|
||||
// La vigilancia apunta a una herramienta del agente y nunca a una URL libre:
|
||||
// es lo que le da el blindaje anti-SSRF sin reimplementarlo.
|
||||
func TestVigilanciaNoAceptaURLLibre(t *testing.T) {
|
||||
b, err := os.ReadFile("umind_vigilancia_tools.go")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
s := string(b)
|
||||
if strings.Contains(s, `"url"`) {
|
||||
t.Error("la tool no puede aceptar una URL: tiene que apuntar a una herramienta ya configurada")
|
||||
}
|
||||
if !strings.Contains(s, "GetUmindHerramientaByNombre(agenteID, herramienta)") {
|
||||
t.Error("falta validar que la herramienta exista y sea de este agente")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,123 @@
|
||||
package services
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
"github.com/sujit-baniya/fiber-boilerplate/pkg/models"
|
||||
)
|
||||
|
||||
// Las tools con las que el dueño arma vigilancias conversando. Solo canal
|
||||
// interno: cada vigilancia es una llamada saliente recurrente y consultas a la
|
||||
// IA, o sea plan del dueño gastado por quien la pida.
|
||||
|
||||
func umindVigilanciaTools(agenteID uint) []agentTool {
|
||||
herramientas, err := models.GetUmindHerramientasActivas(agenteID)
|
||||
if err != nil || len(herramientas) == 0 {
|
||||
return nil // sin herramientas no hay nada que vigilar
|
||||
}
|
||||
nombres := make([]string, len(herramientas))
|
||||
for i, h := range herramientas {
|
||||
nombres[i] = h.Nombre
|
||||
}
|
||||
return []agentTool{
|
||||
{Type: "function", Function: agentToolFunc{
|
||||
Name: "crear_vigilancia",
|
||||
Description: "Programa una vigilancia sobre una API: consulta cada cierto tiempo y avisa cuando se cumple una condición. " +
|
||||
"Herramientas disponibles para vigilar: " + strings.Join(nombres, ", "),
|
||||
Parameters: agentToolParam{
|
||||
Type: "object",
|
||||
Properties: map[string]agentToolParam{
|
||||
"nombre": {Type: "string", Description: "Cómo llamarla, en pocas palabras"},
|
||||
"herramienta": {Type: "string", Description: "Nombre de la herramienta a consultar"},
|
||||
"condicion": {Type: "string", Description: "Qué tiene que pasar para avisar, en lenguaje natural"},
|
||||
"cada_minutos": {Type: "number", Description: "Cada cuántos minutos revisar (mínimo 5, por defecto 30)"},
|
||||
},
|
||||
Required: []string{"nombre", "herramienta", "condicion"},
|
||||
},
|
||||
}},
|
||||
{Type: "function", Function: agentToolFunc{
|
||||
Name: "listar_vigilancias",
|
||||
Description: "Lista las vigilancias activas sobre APIs.",
|
||||
Parameters: agentToolParam{Type: "object", Properties: map[string]agentToolParam{}, Required: []string{}},
|
||||
}},
|
||||
{Type: "function", Function: agentToolFunc{
|
||||
Name: "cancelar_vigilancia",
|
||||
Description: "Apaga una vigilancia. Usá listar_vigilancias para saber el id.",
|
||||
Parameters: agentToolParam{
|
||||
Type: "object",
|
||||
Properties: map[string]agentToolParam{"id": {Type: "number", Description: "El id de la vigilancia"}},
|
||||
Required: []string{"id"},
|
||||
},
|
||||
}},
|
||||
}
|
||||
}
|
||||
|
||||
func executeUmindVigilanciaTool(agenteID uint, sessionID, name string, args map[string]interface{}) string {
|
||||
switch name {
|
||||
case "crear_vigilancia":
|
||||
nombre, _ := args["nombre"].(string)
|
||||
herramienta, _ := args["herramienta"].(string)
|
||||
condicion, _ := args["condicion"].(string)
|
||||
if strings.TrimSpace(nombre) == "" || strings.TrimSpace(condicion) == "" {
|
||||
return `{"error": "falta el nombre o la condición"}`
|
||||
}
|
||||
// La herramienta tiene que existir y ser de este agente: es lo que
|
||||
// hace que la vigilancia no pueda apuntar a una URL arbitraria.
|
||||
if _, err := models.GetUmindHerramientaByNombre(agenteID, herramienta); err != nil {
|
||||
return fmt.Sprintf(`{"error": "no existe una herramienta llamada %q en este agente"}`, herramienta)
|
||||
}
|
||||
cada := 30
|
||||
if n, ok := args["cada_minutos"].(float64); ok && int(n) >= models.UmindVigilanciaIntervaloMin {
|
||||
cada = int(n)
|
||||
}
|
||||
agente, err := models.GetUmindAgenteByID(agenteID)
|
||||
if err != nil {
|
||||
return `{"error": "no se pudo crear la vigilancia"}`
|
||||
}
|
||||
v := &models.UmindVigilancia{
|
||||
AgenteID: agenteID, TenantID: agente.TenantID, Nombre: nombre,
|
||||
Herramienta: herramienta, Condicion: condicion,
|
||||
IntervaloMin: cada, Destino: sessionID, Activa: true,
|
||||
}
|
||||
if err := models.CreateUmindVigilancia(v); err != nil {
|
||||
return fmt.Sprintf(`{"error": "no se pudo crear: ya tenés el máximo de %d vigilancias activas"}`, models.UmindVigilanciaMax)
|
||||
}
|
||||
b, _ := json.Marshal(map[string]interface{}{"ok": true, "id": v.ID, "cada_minutos": cada})
|
||||
return string(b)
|
||||
|
||||
case "listar_vigilancias":
|
||||
items, err := models.GetUmindVigilanciasByAgente(agenteID)
|
||||
if err != nil {
|
||||
return `{"error": "no se pudieron listar"}`
|
||||
}
|
||||
lista := []map[string]interface{}{}
|
||||
for _, v := range items {
|
||||
if !v.Activa {
|
||||
continue
|
||||
}
|
||||
lista = append(lista, map[string]interface{}{
|
||||
"id": v.ID, "nombre": v.Nombre, "condicion": v.Condicion,
|
||||
"cada_minutos": v.IntervaloMin, "en_condicion": v.EnCondicion,
|
||||
})
|
||||
}
|
||||
b, _ := json.Marshal(map[string]interface{}{"vigilancias": lista})
|
||||
return string(b)
|
||||
|
||||
case "cancelar_vigilancia":
|
||||
id, ok := args["id"].(float64)
|
||||
if !ok || id <= 0 {
|
||||
return `{"error": "falta el id"}`
|
||||
}
|
||||
v, err := models.GetUmindVigilanciaByID(uint(id))
|
||||
if err != nil || v.AgenteID != agenteID {
|
||||
return `{"error": "esa vigilancia no existe"}`
|
||||
}
|
||||
if err := models.DeleteUmindVigilancia(uint(id)); err != nil {
|
||||
return `{"error": "no se pudo cancelar"}`
|
||||
}
|
||||
return `{"ok": true}`
|
||||
}
|
||||
return `{"error": "herramienta desconocida"}`
|
||||
}
|
||||
Reference in New Issue
Block a user