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 .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -44,3 +44,4 @@ jobs:
- run: npm run demo:actividad
- run: npm run demo:modo-de-sala
- run: npm run demo:ver-base
- run: npm run demo:candados
1 change: 1 addition & 0 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@
"demo:cripto": "npm run demo:cripto -w server",
"demo:modo-de-sala": "npm run demo:modo-de-sala -w server",
"demo:ver-base": "npm run demo:ver-base -w server",
"demo:candados": "npm run demo:candados -w server",
"demo:modo-persistido": "npm run demo:modo-persistido -w server",
"demo:sala-al-escribir": "npm run demo:sala-al-escribir -w server",
"demo:carga": "npm run demo:carga -w server --",
Expand Down
1 change: 1 addition & 0 deletions server/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@
"demo:politicas": "tsx src/demos/politicas.ts",
"demo:actividad": "tsx src/demos/actividad.ts",
"demo:ver-base": "tsx src/demos/ver-base.ts",
"demo:candados": "tsx src/demos/candados.ts",
"demo:publicar": "tsx src/demos/publicar.ts",
"demo:cuentas": "tsx src/demos/cuentas.ts",
"demo:cripto": "tsx src/demos/cripto.ts",
Expand Down
76 changes: 53 additions & 23 deletions server/src/agent/tools/sql.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
import { mkdir, writeFile } from "node:fs/promises";
import { existsSync } from "node:fs";
import { join } from "node:path";
import { fileMutation } from "../../engine/file-mutation.js";
import { type Tool, ToolError, reqString } from "./base.js";
import { politicasAbiertas, type PoliticaAbierta } from "./politicas.js";

Expand Down Expand Up @@ -77,33 +79,61 @@ export const sqlTool: Tool = {
);
}

const ejecutarSql = ctx.ejecutarSql;
const dir = join(ctx.workspaceDir, "migraciones");
// Un SQL a la vez por sala, con su migración. Dos agentes cambiando la base
// al mismo tiempo (uno crea la tabla, otro la altera) quedaban en un orden
// al azar, y sus migraciones podían caer en el mismo segundo. Es el mismo
// candado de los archivos: la sala ve "esperando a agente-1", y la espera no
// cuenta contra el turno. Y como la migración queda anotada ahí, los demás
// agentes se enteran del cambio en su resumen.
let archivo: string;
try {
await ctx.ejecutarSql(sql);
} catch (err) {
// Lo que diga Postgres (una tabla que ya existe, una columna que no) es
// justo lo que el agente necesita para corregir. Como `error inesperado`
// parecía una falla de Multi, no de su SQL.
const msg = err instanceof Error ? err.message : String(err);
throw new ToolError(
`la base rechazó el SQL y no se guardó migración: ${msg}. Revisa la estructura ` +
"actual con ver_base antes de reintentar.",
archivo = await fileMutation.conCandado(
dir,
{
agentId: ctx.agentId ?? "agente",
onWait: (holder) => ctx.onWaitStart?.({ path: "migraciones", holder }),
},
async () => {
try {
await ejecutarSql(sql);
} catch (err) {
// Lo que diga Postgres (una tabla que ya existe, una columna que no)
// es justo lo que el agente necesita para corregir. Como `error
// inesperado` parecía una falla de Multi, no de su SQL.
const msg = err instanceof Error ? err.message : String(err);
throw new ToolError(
`la base rechazó el SQL y no se guardó migración: ${msg}. Revisa la estructura ` +
"actual con ver_base antes de reintentar.",
);
}

// La migración se guarda DESPUÉS de que corrió, no antes: un archivo
// que describe un cambio que falló es peor que no tenerlo, porque quien
// rehaga la base desde estos archivos acabaría con un esquema que
// nunca existió.
const sello = new Date().toISOString().replace(/[-:T]/g, "").slice(0, 14);
const limpia = descripcion
.toLowerCase()
.replace(/[^a-z0-9]+/g, "-")
.replace(/^-|-$/g, "")
.slice(0, 60);
await mkdir(dir, { recursive: true });
// Dos cambios en el mismo segundo no se pisan: el orden de los archivos
// es el orden en que corrieron.
let nombre = `${sello}-${limpia || "cambio"}.sql`;
for (let n = 2; existsSync(join(dir, nombre)); n++) {
nombre = `${sello}-${n}-${limpia || "cambio"}.sql`;
}
await writeFile(join(dir, nombre), sql.trim() + "\n", "utf8");
return { valor: nombre, tocados: [join(dir, nombre)] };
},
);
} finally {
ctx.onWaitEnd?.();
}

// La migración se guarda DESPUÉS de que corrió, no antes: un archivo que
// describe un cambio que falló es peor que no tenerlo, porque quien rehaga
// la base desde estos archivos acabaría con un esquema que nunca existió.
const sello = new Date().toISOString().replace(/[-:T]/g, "").slice(0, 14);
const limpia = descripcion
.toLowerCase()
.replace(/[^a-z0-9]+/g, "-")
.replace(/^-|-$/g, "")
.slice(0, 60);
const dir = join(ctx.workspaceDir, "migraciones");
await mkdir(dir, { recursive: true });
const archivo = `${sello}-${limpia || "cambio"}.sql`;
await writeFile(join(dir, archivo), sql.trim() + "\n", "utf8");

ctx.emit?.({ type: "file:changed", path: `migraciones/${archivo}`, action: "write" });

const listo = `listo. El cambio quedó en migraciones/${archivo}`;
Expand Down
175 changes: 175 additions & 0 deletions server/src/demos/candados.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,175 @@
import { mkdtemp, readdir, rm } from "node:fs/promises";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { sqlTool } from "../agent/tools/sql.js";
import type { ToolContext } from "../agent/tools/base.js";
import { fileMutation } from "../engine/file-mutation.js";
import { AgentRegistry, resumenDeOtros } from "../engine/agents.js";
import { SqliteStorage } from "../storage/sqlite.js";
import { startTurn, listTurns } from "../engine/turns.js";

/**
* Demo: lo que no puede pasar dos veces a la vez, no pasa.
* Uso: npm run demo:candados
*
* 1. La tool sql corre de a uno por sala, y lo que hizo llega al resumen de los
* demás agentes, sin mezclar salas.
* 2. La conexión de Supabase se actualiza por campos: renovar el token no borra
* el proyecto, y guardar el proyecto no regresa tokens viejos.
*
* No necesita red: el SQL "corre" contra una función que tarda y anota, y la
* conexión se guarda en una base SQLite temporal.
*/

let pass = 0;
let fail = 0;
function check(name: string, ok: boolean, detail = ""): void {
if (ok) {
pass++;
console.log(` [ok] ${name}`);
} else {
fail++;
console.log(` [X] ${name} ${detail}`);
}
}

const dormir = (ms: number) => new Promise((r) => setTimeout(r, ms));

/** Un ejecutarSql que tarda y cuenta cuántos corren a la vez. */
function baseLenta() {
const estado = { ahora: 0, maximo: 0, orden: [] as string[] };
const ejecutar = async (sql: string) => {
estado.ahora++;
estado.maximo = Math.max(estado.maximo, estado.ahora);
await dormir(150);
estado.orden.push(sql);
estado.ahora--;
};
return { estado, ejecutar };
}

async function main() {
console.log("\n=== candados ===\n");
const raiz = await mkdtemp(join(tmpdir(), "multi-candados-"));

console.log("1. sql corre de a uno por sala");
{
const sala = join(raiz, "sala-a");
const { estado, ejecutar } = baseLenta();
const esperas: Array<{ path: string; holder?: string }> = [];
const ctx = (agentId: string) =>
({
workspaceDir: sala,
agentId,
ejecutarSql: ejecutar,
onWaitStart: (i: { path: string; holder?: string }) => esperas.push(i),
onWaitEnd: () => {},
}) as unknown as ToolContext;

await Promise.all([
sqlTool.run({ sql: "create table tareas (id bigint)", descripcion: "crea-tabla-tareas" }, ctx("agente-1")),
sqlTool.run({ sql: "alter table tareas add column hecha boolean", descripcion: "agrega-hecha" }, ctx("agente-2")),
]);
check("nunca corrieron dos a la vez", estado.maximo === 1, `máximo ${estado.maximo}`);
check("en el orden en que llegaron", estado.orden[0].startsWith("create table"), estado.orden.join(" | "));
check(
"el segundo supo a quién esperaba",
esperas.length === 1 && esperas[0].holder === "agente-1" && esperas[0].path === "migraciones",
JSON.stringify(esperas),
);
const archivos = (await readdir(join(sala, "migraciones"))).sort();
check("dos migraciones, sin pisarse aunque caigan en el mismo segundo", archivos.length === 2, archivos.join(", "));
}

console.log("\n2. Dos salas no se esperan entre sí");
{
const { estado, ejecutar } = baseLenta();
const ctx = (sala: string) =>
({ workspaceDir: join(raiz, sala), agentId: "agente-1", ejecutarSql: ejecutar }) as unknown as ToolContext;
await Promise.all([
sqlTool.run({ sql: "create table a (id bigint)", descripcion: "a" }, ctx("sala-b")),
sqlTool.run({ sql: "create table b (id bigint)", descripcion: "b" }, ctx("sala-c")),
]);
check("corrieron al mismo tiempo", estado.maximo === 2, `máximo ${estado.maximo}`);
}

console.log("\n3. Los demás se enteran, y solo los de la misma sala");
{
const deA = fileMutation.trabajoRecienteDeOtros("agente-9", { sala: join(raiz, "sala-a") });
check(
"en la sala A se ven sus dos migraciones",
deA.filter((f) => f.path.startsWith("migraciones/")).length === 2,
JSON.stringify(deA.map((f) => f.path)),
);
check("y no las de las salas B y C, aunque sean de un agente-1", deA.every((f) => !f.path.includes("-a.sql") && !f.path.includes("-b.sql")));
check("con rutas relativas a la sala", deA.every((f) => !f.path.startsWith("/")), deA[0]?.path);

const reg = new AgentRegistry();
const a1 = reg.spawn("haz las tareas")!;
const a2 = reg.spawn("agrega un filtro")!;
reg.finish(a1.id);
const resumen = resumenDeOtros(reg, a2.id, fileMutation.trabajoRecienteDeOtros(a2.id, { sala: join(raiz, "sala-a") })) ?? "";
check("el resumen dice qué cambió en la base", resumen.includes("cambió la base: crea-tabla-tareas"), resumen);
check("sin listar la migración como archivo por releer", !resumen.includes("ya tocó: migraciones"), resumen);
check("y le dice que revise con ver_base", resumen.includes("ver_base"));
}

console.log("\n4. La conexión de Supabase se actualiza por campos");
{
process.env.MULTI_LLAVE ??= "llave-de-la-demo";
const storage = new SqliteStorage(join(raiz, "multi.db"));
await storage.init();
await storage.createRoom({ id: "sala-sb", workspaceDir: join(raiz, "sala-sb"), createdAt: Date.now(), lastActiveAt: Date.now() });
await storage.guardarConexionSupabase({
roomId: "sala-sb",
acceso: "acceso-1",
refresco: "refresco-1",
expiraEn: 1,
conectadoEn: Date.now(),
});

// La preparación guarda el proyecto; en medio, alguien renovó el token.
await storage.actualizarConexionSupabase("sala-sb", { acceso: "acceso-2", refresco: "refresco-2", expiraEn: 2 });
await storage.actualizarConexionSupabase("sala-sb", { proyecto: "refdelproyecto", password: "secreta" });
let c = await storage.conexionSupabase("sala-sb");
check("guardar el proyecto no regresa los tokens viejos", c?.acceso === "acceso-2" && c?.refresco === "refresco-2", JSON.stringify(c));
check("y el proyecto quedó", c?.proyecto === "refdelproyecto" && c?.password === "secreta");

// Otra renovación después: no toca el proyecto.
await storage.actualizarConexionSupabase("sala-sb", { acceso: "acceso-3", refresco: "refresco-3", expiraEn: 3 });
c = await storage.conexionSupabase("sala-sb");
check("renovar el token no borra el proyecto ni la contraseña", c?.proyecto === "refdelproyecto" && c?.password === "secreta", JSON.stringify(c));
check("y los tokens son los nuevos", c?.acceso === "acceso-3" && c?.expiraEn === 3);

await storage.actualizarConexionSupabase("sin-conexion", { acceso: "x" });
check("actualizar una sala sin conexión no crea nada", (await storage.conexionSupabase("sin-conexion")) === null);
}

console.log("\n5. Dos turnos que empiezan a la vez no tumban el server");
{
// Antes compartían el nombre temporal: el segundo rename tronaba con ENOENT
// sin que nadie lo atrapara, y se caía el server con todas las salas. Y el
// que sí escribía borraba el turno del otro.
const sala = join(raiz, "sala-turnos");
let error = "";
try {
await Promise.all(
Array.from({ length: 5 }, (_, i) => startTurn(sala, { roomId: "sala-turnos", agentId: `agente-${i + 1}`, task: `tarea ${i}` })),
);
} catch (err) {
error = String(err);
}
check("cinco a la vez, sin error", error === "", error);
const turnos = await listTurns(sala);
check("y quedan los cinco", turnos.length === 5, `quedaron ${turnos.length}`);
}

await rm(raiz, { recursive: true, force: true });
console.log(`\n${pass} pasaron, ${fail} fallaron\n`);
process.exit(fail > 0 ? 1 : 0);
}

main().catch((err) => {
console.error("\ndemo falló:", err);
process.exit(1);
});
20 changes: 18 additions & 2 deletions server/src/engine/agents.ts
Original file line number Diff line number Diff line change
Expand Up @@ -241,7 +241,11 @@ export function resumenDeOtros(
if (otros.length === 0) return null;

const lineas = otros.map((a) => {
const suyos = archivos.filter((f) => f.agentId === a.id);
const todos = archivos.filter((f) => f.agentId === a.id);
// Las migraciones son cambios a la BASE, no archivos que haya que releer: se
// cuentan aparte y por lo que hicieron ("crea-tabla-tareas").
const base = todos.filter((f) => esMigracion(f.path)).map((f) => cambioDeBase(f.path));
const suyos = todos.filter((f) => !esMigracion(f.path));
const enVuelo = suyos.filter((f) => f.escribiendoAhora).map((f) => f.path);
const tocados = suyos.filter((f) => !f.escribiendoAhora).map((f) => f.path);

Expand All @@ -251,6 +255,7 @@ export function resumenDeOtros(
const partes = [`${a.name}${estado}: ${a.task ?? "trabajando"}`];
if (enVuelo.length) partes.push(` escribiendo ahora: ${enVuelo.join(", ")}`);
if (tocados.length) partes.push(` ya tocó: ${tocados.slice(0, 8).join(", ")}`);
if (base.length) partes.push(` cambió la base: ${base.slice(0, 8).join(", ")}`);

const dijo = ultimoMensaje?.get(a.id);
if (dijo) partes.push(` dijo: "${recorta(dijo, 300)}"`);
Expand All @@ -266,14 +271,25 @@ export function resumenDeOtros(
"No rehagas lo que otro ya hizo ni lo que está haciendo. Si alguien está montando",
"el proyecto, espera a que termine o trabaja en otra parte. Un archivo que otro",
"está escribiendo en este momento: déjalo. Uno que ya soltó: léelo antes de",
"tocarlo, porque cambió desde la última vez que lo viste.",
"tocarlo, porque cambió desde la última vez que lo viste. Si otro cambió la base,",
"revisa con ver_base antes de tocar esas tablas.",
"",
"Esto cubre los últimos minutos. Si lo que te piden pudo hacerse antes, `git log`",
"dice qué hizo cada agente en cada turno.",
"</otros_agentes>",
].join("\n");
}

function esMigracion(path: string): boolean {
return /(^|[\\/])migraciones[\\/][^\\/]+\.sql$/.test(path);
}

/** `migraciones/20260924065501-crea-tabla-tareas.sql` → `crea-tabla-tareas`. */
function cambioDeBase(path: string): string {
const nombre = path.split(/[\\/]/).pop() ?? path;
return nombre.replace(/\.sql$/, "").replace(/^\d{14}-(\d+-)?/, "");
}

/** Recorta a lo que quepa sin comerse el contexto del turno. */
function recorta(s: string, max: number): string {
const limpio = s.replace(/\s+/g, " ").trim();
Expand Down
Loading
Loading