From e24241ead9cec41d968794aad6f3716f91760393 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 24 Sep 2026 17:26:23 +0000 Subject: [PATCH] Candados: la conexion de Supabase, la tool sql y el archivo de turnos Conexion de Supabase: la vuelta de autorizar, la preparacion y la renovacion del token guardaban la fila completa desde una copia leida antes. - Una segunda autorizacion (dos pestanas, doble clic) borraba el proyecto y la contrasena; si era de otra cuenta, la sala acababa en una base nueva y vacia. Ahora con proyecto solo se renuevan los tokens si esa cuenta lo ve, y si no, no se toca nada y el panel dice por que. - actualizarConexionSupabase cambia solo los campos que llegan. - La renovacion del token va de una en una por sala y relee adentro. Tool sql: un SQL a la vez por sala con su migracion, con el mismo candado de los archivos (la sala ve "esperando a agente-1"). La migracion queda anotada, asi que los demas agentes ven "cambio la base: crea-tabla-x". Dos migraciones en el mismo segundo ya no se pisan. Resumen de otros agentes: filtraba solo por id de agente, y los ids se repiten en cada sala, asi que llegaban archivos de otras salas. Ahora filtra por la carpeta de la sala y da rutas relativas. Turnos: dos agentes que empezaban turno a la vez compartian el archivo temporal; el segundo rename tronaba sin que nadie lo atrapara y tumbaba el server. Ahora las escrituras van en fila por sala y el temporal es unico. demo:candados cubre todo esto y entra al CI. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01HJ5oT2Xm3VZMvbPAcgKbz4 --- .github/workflows/ci.yml | 1 + package.json | 1 + server/package.json | 1 + server/src/agent/tools/sql.ts | 76 +++++++++---- server/src/demos/candados.ts | 175 +++++++++++++++++++++++++++++ server/src/engine/agents.ts | 20 +++- server/src/engine/file-mutation.ts | 43 ++++++- server/src/engine/turns.ts | 63 +++++++---- server/src/index.ts | 153 ++++++++++++++++++------- server/src/storage/sqlite.ts | 32 ++++++ server/src/storage/types.ts | 12 ++ 11 files changed, 490 insertions(+), 87 deletions(-) create mode 100644 server/src/demos/candados.ts diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 092cb7e..0e9c762 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -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 diff --git a/package.json b/package.json index 79d1389..646519f 100644 --- a/package.json +++ b/package.json @@ -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 --", diff --git a/server/package.json b/server/package.json index b0561b1..0bff97e 100644 --- a/server/package.json +++ b/server/package.json @@ -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", diff --git a/server/src/agent/tools/sql.ts b/server/src/agent/tools/sql.ts index 9d8dc99..daa23a4 100644 --- a/server/src/agent/tools/sql.ts +++ b/server/src/agent/tools/sql.ts @@ -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"; @@ -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}`; diff --git a/server/src/demos/candados.ts b/server/src/demos/candados.ts new file mode 100644 index 0000000..ea5e376 --- /dev/null +++ b/server/src/demos/candados.ts @@ -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); +}); diff --git a/server/src/engine/agents.ts b/server/src/engine/agents.ts index e543673..da8c485 100644 --- a/server/src/engine/agents.ts +++ b/server/src/engine/agents.ts @@ -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); @@ -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)}"`); @@ -266,7 +271,8 @@ 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.", @@ -274,6 +280,16 @@ export function resumenDeOtros( ].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(); diff --git a/server/src/engine/file-mutation.ts b/server/src/engine/file-mutation.ts index b43c384..975173a 100644 --- a/server/src/engine/file-mutation.ts +++ b/server/src/engine/file-mutation.ts @@ -1,6 +1,6 @@ import { readFile, writeFile, rename, mkdir } from "node:fs/promises"; import { existsSync } from "node:fs"; -import { dirname, resolve } from "node:path"; +import { dirname, relative, resolve, sep } from "node:path"; import { KeyedMutex } from "./keyed-mutex.js"; /** @@ -113,6 +113,31 @@ export class FileMutation { ); } + /** + * Corre `fn` con el candado de `ruta` tomado, igual que una escritura. + * + * Para lo que no es escribir un archivo pero tampoco puede pasar dos veces a + * la vez: correr SQL contra la base de la sala, por ejemplo. Lo que `fn` + * devuelva en `tocados` queda anotado como escrito por `agentId`, así que + * entra al resumen que reciben los demás agentes. + */ + async conCandado( + ruta: string, + opts: { agentId: string; onWait?: (holder: string | undefined) => void }, + fn: () => Promise<{ valor: T; tocados?: string[] }>, + ): Promise { + return this.mutex.run( + resolve(ruta), + async () => { + const { valor, tocados = [] } = await fn(); + for (const t of tocados) this.lastWriter.set(resolve(t), { agentId: opts.agentId, at: Date.now() }); + this.pruneWriters(); + return valor; + }, + { owner: opts.agentId, onWait: opts.onWait }, + ); + } + /** Lee un archivo (fuera del lock: leer no necesita exclusión). */ async read(path: string): Promise { const key = resolve(path); @@ -133,12 +158,24 @@ export class FileMutation { */ trabajoRecienteDeOtros( exceptoAgente: string, - dentroDeMs = 120_000, + opts: { + dentroDeMs?: number; + /** + * La carpeta de la sala. Sin esto se mezclaban salas: los ids de agente se + * repiten en cada una (`agente-1`), y este registro es uno por proceso, así + * que el agente-2 de una sala recibía lo que tocó el agente-1 de otra. Con + * ella, además, las rutas salen relativas a la sala y no absolutas. + */ + sala?: string; + } = {}, ): Array<{ agentId: string; path: string; hace: number; escribiendoAhora: boolean }> { + const { dentroDeMs = 120_000 } = opts; + const sala = opts.sala ? resolve(opts.sala) + sep : null; const ahora = Date.now(); const out: Array<{ agentId: string; path: string; hace: number; escribiendoAhora: boolean }> = []; for (const [path, w] of this.lastWriter) { if (w.agentId === exceptoAgente) continue; + if (sala && !path.startsWith(sala)) continue; const hace = ahora - w.at; if (hace > dentroDeMs) continue; // Con el lock tomado no es "lo tocó": lo está escribiendo AHORA. Distinguirlo @@ -146,7 +183,7 @@ export class FileMutation { // hay que releerlo. out.push({ agentId: w.agentId, - path, + path: sala ? relative(sala, path) : path, hace, escribiendoAhora: !!this.mutex.info(path).holder, }); diff --git a/server/src/engine/turns.ts b/server/src/engine/turns.ts index 4a5fe45..ed41929 100644 --- a/server/src/engine/turns.ts +++ b/server/src/engine/turns.ts @@ -2,6 +2,7 @@ import { readFile, writeFile, mkdir, rename } from "node:fs/promises"; import { existsSync } from "node:fs"; import { join, dirname } from "node:path"; import { commitAll } from "./git.js"; +import { KeyedMutex } from "./keyed-mutex.js"; /** * Turnos con estado DURABLE (en disco, no memoria). @@ -53,12 +54,37 @@ async function readTurns(workspaceDir: string): Promise { async function writeTurns(workspaceDir: string, turns: Turn[]): Promise { const f = turnsFile(workspaceDir); await mkdir(dirname(f), { recursive: true }); - // Atómico: si el server muere escribiendo, no queda un JSON corrupto. - const tmp = `${f}.tmp-${process.pid}`; + // Atómico: si el server muere escribiendo, no queda un JSON corrupto. El + // nombre temporal es único por escritura: con uno por proceso, dos escrituras + // a la vez renombraban el mismo archivo y la segunda tronaba con ENOENT. + const tmp = `${f}.tmp-${process.pid}-${Date.now()}-${Math.random().toString(36).slice(2, 8)}`; await writeFile(tmp, JSON.stringify(turns, null, 2), "utf8"); await rename(tmp, f); } +/** + * Leer, cambiar y escribir el archivo de turnos de una sala, de uno en uno. + * + * Dos agentes que empiezan turno a la vez (dos personas escriben @agente en el + * mismo momento) leían los dos la lista, agregaban cada uno el suyo y el último + * en escribir borraba el turno del otro. Y como compartían el nombre temporal, + * el segundo `rename` tronaba sin que nadie lo atrapara y TUMBABA EL SERVER, + * con todas las salas. + */ +const candadoDeTurnos = new KeyedMutex(); + +async function modificarTurnos(workspaceDir: string, fn: (turns: Turn[]) => T): Promise { + return candadoDeTurnos.run(turnsFile(workspaceDir), async () => { + const turns = await readTurns(workspaceDir); + const antes = JSON.stringify(turns); + const resultado = fn(turns); + // Sin cambios no se escribe: el barrido de arranque pasa por todas las + // salas, y no tiene por qué crearle un archivo a cada una. + if (JSON.stringify(turns) !== antes) await writeTurns(workspaceDir, turns); + return resultado; + }); +} + /** Abre un turno (estado "running" en disco). */ export async function startTurn( workspaceDir: string, @@ -73,9 +99,9 @@ export async function startTurn( startedAt: Date.now(), activeMs: 0, }; - const turns = await readTurns(workspaceDir); - turns.push(turn); - await writeTurns(workspaceDir, turns); + await modificarTurnos(workspaceDir, (turns) => { + turns.push(turn); + }); return turn; } @@ -129,11 +155,10 @@ export async function failTurnConCommit( } async function patchTurn(workspaceDir: string, turnId: string, patch: Partial): Promise { - const turns = await readTurns(workspaceDir); - const t = turns.find((x) => x.id === turnId); - if (!t) return; - Object.assign(t, patch); - await writeTurns(workspaceDir, turns); + await modificarTurnos(workspaceDir, (turns) => { + const t = turns.find((x) => x.id === turnId); + if (t) Object.assign(t, patch); + }); } /** @@ -145,16 +170,14 @@ async function patchTurn(workspaceDir: string, turnId: string, patch: Partial { - const turns = await readTurns(workspaceDir); - const orphans = turns.filter((t) => t.state === "running"); - if (orphans.length === 0) return []; - - for (const t of orphans) { - t.state = "orphaned"; - t.endedAt = Date.now(); - } - await writeTurns(workspaceDir, turns); - return orphans; + return modificarTurnos(workspaceDir, (turns) => { + const orphans = turns.filter((t) => t.state === "running"); + for (const t of orphans) { + t.state = "orphaned"; + t.endedAt = Date.now(); + } + return orphans; + }); } export async function listTurns(workspaceDir: string): Promise { diff --git a/server/src/index.ts b/server/src/index.ts index 7f860a4..34922f7 100644 --- a/server/src/index.ts +++ b/server/src/index.ts @@ -800,6 +800,9 @@ fastify.get<{ Params: { id: string } }>("/rooms/:id/supabase", async (req) => { proyecto: conexion?.proyecto ?? null, url: conexion?.proyecto ? urlDelProyecto(conexion.proyecto) : null, etapa: etapaSupabase.get(req.params.id) ?? null, + // Lo que salió mal al volver de autorizar. Va aquí y no solo por socket + // porque esa vuelta RECARGA la página: el aviso llegaría antes que la Sala. + error: errorSupabase.get(req.params.id) ?? null, // Tener proyecto no es lo mismo que tener las variables en el `.env`, que // es lo único que el agente ve. El panel decía "conectada" en cuanto el // proyecto existía, y si la preparación se cortaba después (el tope de @@ -884,24 +887,24 @@ fastify.get<{ Querystring: { code?: string; state?: string; error?: string } }>( if (!tokens) return reply.code(502).send({ error: "Supabase no confirmó la autorización" }); console.log(`[supabase] tokens recibidos, redirigiendo a ${aLaSala}`); - // Se guarda ANTES de crear el proyecto, a propósito: crear tarda minutos y - // si el server se cae en medio, la autorización no se pierde y se puede - // reintentar solo la parte que faltó. - await (await getStorage()).guardarConexionSupabase({ - roomId: vuelo.roomId, - acceso: tokens.acceso, - refresco: tokens.refresco, - expiraEn: tokens.expiraEn, - conectadoEn: Date.now(), - }); - - // La persona se va de vuelta a su sala YA. Crear el proyecto sigue por su - // cuenta y la sala se entera por socket: tenerla mirando una pantalla de - // redirección durante tres minutos sería peor. - void prepararProyecto(vuelo.roomId).catch((err) => { + const fallo = (err: unknown) => { + const msg = err instanceof Error ? err.message : String(err); console.error(`[supabase] no se pudo preparar el proyecto de ${vuelo.roomId}:`, err); - io.to(vuelo.roomId).emit("supabase:fallo", { error: String(err?.message ?? err) }); - }); + io.to(vuelo.roomId).emit("supabase:fallo", { error: msg }); + }; + // Guardar va en la misma fila que la preparación de esta sala, para que una + // segunda autorización no se meta a media creación. Crear el proyecto sigue + // por su cuenta y la sala se entera por socket: tenerla mirando una pantalla + // de redirección durante minutos sería peor. + const guardado = preparacionesDeSupabase + .run(vuelo.roomId, () => guardarAutorizacion(vuelo.roomId, tokens)) + .then((preparar) => { + if (preparar) void prepararProyecto(vuelo.roomId).catch(fallo); + }, fallo); + // Se espera a que quede guardado antes de volver, como antes: así el panel + // ya sabe al recargar que la base se está preparando (o por qué no). Con + // tope, porque si otra preparación de la sala va en curso, eso son minutos. + await Promise.race([guardado, new Promise((r) => setTimeout(r, 3000))]); return reply.redirect(aLaSala); }, @@ -926,6 +929,67 @@ const preparacionesDeSupabase = new KeyedMutex(); */ const etapaSupabase = new Map(); +/** + * Lo que salió mal en la última vuelta de autorizar, por sala. Lo lee el `GET` + * del panel; se borra con la siguiente autorización buena o al desconectar. + */ +const errorSupabase = new Map(); + +/** + * Guarda lo que volvió de autorizar sin pisar lo que la sala ya tenía. Devuelve + * si hay que preparar el proyecto. + * + * Antes se guardaba la fila completa con proyecto y contraseña en null. Una + * segunda autorización (dos pestañas, un doble clic, dos personas de la sala a + * la vez) borraba así el proyecto y la contraseña, que Supabase no devuelve + * nunca; y si venía de OTRA cuenta, la sala acababa con una base nueva y vacía. + * + * Corre dentro de `preparacionesDeSupabase`: no se mete a media creación. + */ +async function guardarAutorizacion( + roomId: string, + tokens: { acceso: string; refresco: string; expiraEn: number }, +): Promise { + const storage = await getStorage(); + const previa = await storage.conexionSupabase(roomId); + + if (!previa) { + // Se guarda ANTES de crear el proyecto, a propósito: crear tarda minutos y + // si el server se cae en medio, la autorización no se pierde y se puede + // reintentar solo la parte que faltó. + await storage.guardarConexionSupabase({ roomId, ...tokens, conectadoEn: Date.now() }); + errorSupabase.delete(roomId); + return true; + } + + if (previa.proyecto) { + // Con base ya hecha, solo cuenta si es la MISMA cuenta: la que puede ver ese + // proyecto. Otra cuenta no se lleva la sala; para cambiarla hay que + // desconectar primero, que es una decisión y no un accidente. + const suyos = await proyectos(tokens.acceso); + const esLaMisma = suyos.some((p) => (p.ref ?? p.id) === previa.proyecto); + if (!esLaMisma) { + console.log(`[supabase] ${roomId} autorizó otra cuenta; la conexión no se cambia`); + const error = + "Esa cuenta de Supabase no ve la base de esta sala (es otra cuenta, o el proyecto " + + "ya no existe), así que no se cambió nada. Para usar otra, desconecta primero."; + errorSupabase.set(roomId, error); + io.to(roomId).emit("supabase:fallo", { error }); + return false; + } + await storage.actualizarConexionSupabase(roomId, tokens); + errorSupabase.delete(roomId); + console.log(`[supabase] ${roomId} volvió a autorizar la misma cuenta: solo se renuevan los tokens`); + return false; + } + + // Autorizada y sin proyecto todavía: la preparación quedó a medias o está en + // fila. Tokens frescos, y que siga. + await storage.actualizarConexionSupabase(roomId, tokens); + errorSupabase.delete(roomId); + return true; +} + async function prepararProyecto(roomId: string): Promise { // Serializado por sala: dos caminos llaman aquí (el callback al autorizar, y // el botón de conectar cuando retoma), y si coinciden los dos leen que no hay @@ -985,7 +1049,7 @@ async function prepararProyectoSerializado(roomId: string): Promise { console.log(`[supabase] ${roomId} ya tenía el proyecto ${ref}: se adopta en vez de crear otro`); // La contraseña no se recupera: Supabase no la devuelve nunca. La app no // la necesita; solo quien quiera entrar a la base por fuera. - await (await getStorage()).guardarConexionSupabase({ ...conexion, proyecto: ref, password: null }); + await (await getStorage()).actualizarConexionSupabase(roomId, { proyecto: ref, password: null }); } } @@ -1013,7 +1077,7 @@ async function prepararProyectoSerializado(roomId: string): Promise { // Se guarda la referencia en cuanto existe, aunque el proyecto todavía esté // levantándose: si el server se reinicia ahora, lo que ya se creó no se // vuelve a crear. - await (await getStorage()).guardarConexionSupabase({ ...conexion, proyecto: ref, password }); + await (await getStorage()).actualizarConexionSupabase(roomId, { proyecto: ref, password }); } else { console.log(`[supabase] retomando el proyecto ${ref} de ${roomId}`); } @@ -1232,35 +1296,46 @@ async function conLaBaseDeLaSala( * hay nada que reintentar, hay que volver a autorizar. */ async function accesoVigente(roomId: string): Promise { - const storage = await getStorage(); - const conexion = await storage.conexionSupabase(roomId); - if (!conexion) return null; + // Una renovación a la vez por sala. Dos al mismo tiempo (dos agentes, o sql y + // ver_base) mandaban el MISMO token de renovación, y Supabase los rota: el + // segundo fallaba y la sala se quedaba "sin poder renovar". Adentro se relee, + // así el que esperó usa lo que el primero acaba de renovar. + return renovacionesDeSupabase.run(roomId, async () => { + const storage = await getStorage(); + const conexion = await storage.conexionSupabase(roomId); + if (!conexion) return null; - // Un minuto de margen: si está a punto de caducar, mejor renovar ahora que a - // media operación. - if (conexion.expiraEn > Date.now() + 60_000) return conexion.acceso; + // Un minuto de margen: si está a punto de caducar, mejor renovar ahora que a + // media operación. + if (conexion.expiraEn > Date.now() + 60_000) return conexion.acceso; - const cred = credencialDeSupabase(); - if (!cred) return null; - const tokens = await refrescarSupabase(cred, conexion.refresco); - if (!tokens) { - console.error(`[supabase] ${roomId} ya no puede renovar su token`); - return null; - } - await storage.guardarConexionSupabase({ - ...conexion, - acceso: tokens.acceso, - refresco: tokens.refresco, - expiraEn: tokens.expiraEn, + const cred = credencialDeSupabase(); + if (!cred) return null; + const tokens = await refrescarSupabase(cred, conexion.refresco); + if (!tokens) { + console.error(`[supabase] ${roomId} ya no puede renovar su token`); + return null; + } + // Solo los tokens: guardar la fila completa desde la copia de arriba podía + // regresar el proyecto a null si la preparación lo guardó mientras tanto. + await storage.actualizarConexionSupabase(roomId, { + acceso: tokens.acceso, + refresco: tokens.refresco, + expiraEn: tokens.expiraEn, + }); + return tokens.acceso; }); - return tokens.acceso; } +/** Ver `accesoVigente`. Aparte del de preparaciones: la preparación lo llama adentro. */ +const renovacionesDeSupabase = new KeyedMutex(); + fastify.delete<{ Params: { id: string } }>("/rooms/:id/supabase", async (req) => { // Solo se desconecta de aquí: el proyecto en Supabase sigue existiendo porque // es de la persona, y las variables se quedan en el .env porque la app las // sigue necesitando para funcionar. await (await getStorage()).borrarConexionSupabase(req.params.id); + errorSupabase.delete(req.params.id); // Si vuelve a conectar, puede ser otro proyecto: hay que revisarlo de nuevo. loginAnonimoRevisado.delete(req.params.id); io.to(req.params.id).emit("supabase:desconectado", {}); @@ -2252,7 +2327,7 @@ async function runAgentTurn( * algo variable tiraría ese caché en cada llamada. */ function conContextoDeOtros(room: Room, agentId: string, task: string): string { - const archivos = fileMutation.trabajoRecienteDeOtros(agentId); + const archivos = fileMutation.trabajoRecienteDeOtros(agentId, { sala: room.workspace.dir }); const resumen = resumenDeOtros(room.agents, agentId, archivos, ultimosMensajes(room)); return resumen ? `${resumen}\n\n${task}` : task; } diff --git a/server/src/storage/sqlite.ts b/server/src/storage/sqlite.ts index 3765320..455bbdb 100644 --- a/server/src/storage/sqlite.ts +++ b/server/src/storage/sqlite.ts @@ -486,6 +486,38 @@ export class SqliteStorage implements Storage { return rows.map((r) => r.room_id); } + async actualizarConexionSupabase( + roomId: string, + campos: Partial>, + ): Promise { + const columnas: string[] = []; + const valores: (string | number | null)[] = []; + if (campos.acceso !== undefined) { + columnas.push("acceso = ?"); + valores.push(cifrar(campos.acceso)); + } + if (campos.refresco !== undefined) { + columnas.push("refresco = ?"); + valores.push(cifrar(campos.refresco)); + } + if (campos.expiraEn !== undefined) { + columnas.push("expira_en = ?"); + valores.push(campos.expiraEn); + } + if (campos.proyecto !== undefined) { + columnas.push("proyecto = ?"); + valores.push(campos.proyecto ?? null); + } + if (campos.password !== undefined) { + columnas.push("password = ?"); + valores.push(campos.password ? cifrar(campos.password) : null); + } + if (columnas.length === 0) return; + this.db + .prepare(`UPDATE supabase_salas SET ${columnas.join(", ")} WHERE room_id = ?`) + .run(...valores, roomId); + } + async borrarConexionSupabase(roomId: string): Promise { this.db.prepare(`DELETE FROM supabase_salas WHERE room_id = ?`).run(roomId); } diff --git a/server/src/storage/types.ts b/server/src/storage/types.ts index 35be56f..a39246a 100644 --- a/server/src/storage/types.ts +++ b/server/src/storage/types.ts @@ -221,6 +221,18 @@ export interface Storage { */ conexionSupabase(roomId: string): Promise; guardarConexionSupabase(conexion: ConexionSupabase): Promise; + /** + * Cambia SOLO los campos que llegan, sin tocar los demás. + * + * Es lo que usan la renovación del token, la preparación y una segunda + * autorización. Guardar la fila completa desde una copia leída antes pisaba lo + * que otro camino había escrito mientras tanto: el proyecto volvía a null, o + * los tokens a unos ya rotados. No hace nada si la sala no tiene conexión. + */ + actualizarConexionSupabase( + roomId: string, + campos: Partial>, + ): Promise; borrarConexionSupabase(roomId: string): Promise; /** Las salas que tienen una conexión guardada, para retomar las que quedaron a medias. */ salasConSupabase(): Promise;