En este tutorial, creamos un flujo de trabajo de IA agente ultraavanzado que se comporta como un sistema de investigación y razonamiento de nivel de producción en lugar de una única llamada rápida. Ingerimos fuentes web reales de forma asincrónica, las dividimos en fragmentos con seguimiento de procedencia y ejecutamos una recuperación híbrida utilizando incrustaciones TF-IDF (dispersas) y OpenAI (densas), luego fusionamos los resultados para una mayor recuperación y estabilidad. Orquestamos múltiples agentes, planificación, síntesis y reparación, mientras aplicamos barreras estrictas para que cada reclamo importante se base en evidencia recuperada y persistamos en la memoria episódica. Por tanto, el sistema mejora su estrategia con el tiempo. Consulta los CÓDIGOS COMPLETOS aquí.
!pip -q instalar openai openai-agents pydantic httpx beautifulsoup4 lxml scikit-learn numpy import os, re, json, time, getpass, asyncio, sqlite3, hashlib escribiendo import List, Dict, Tuple, Opcional, Cualquier importación numpy como np import httpx desde bs4 import BeautifulSoup desde pydantic import BaseModel, Field from sklearn.feature_extraction.text importar TfidfVectorizer desde sklearn.metrics.pairwise importar cosine_similarity desde openai importar AsyncOpenAI desde agentes importar Agent, Runner, SQLiteSession si no os.environ.get("OPENAI_API_KEY"): os.environ["OPENAI_API_KEY"] = getpass.getpass("Ingrese su clave API de OpenAI: ") si no es os.environ.get("OPENAI_API_KEY"): elevar RuntimeError("OPENAI_API_KEY no proporcionada.") print("✅ Clave API de OpenAI cargada de forma segura.") oa = AsyncOpenAI(api_key=os.environ["OPENAI_API_KEY"]) def sha1(s: str) -> str: return hashlib.sha1(s.encode("utf-8", errores="ignore")).hexdigest() def normalize_url(u: str) -> str: u = (u o "").strip() return u.rstrip()".,]"'") def clean_html_to_text(html: str) -> str: sopa = BeautifulSoup(html, "lxml") para etiqueta en sopa(["script", "style", "noscript"]): tag.decompose() txt = sopa.get_text("n") txt = re.sub(r"n{3,}", "nn", txt).strip() txt = re.sub(r"[ t]+", " ", txt) return txt def chunk_text(texto: str, chunk_chars: int = 1600, solapamiento_chars: int = 320) -> Lista[cadena]: si no texto: regresar[]texto = re.sub(r"s+", " ", texto).strip() n = len(texto) paso = max(1, fragmentos_caracteres – superposición_caracteres) fragmentos =[]i = 0 mientras i < n: trozos.append(text[i:i + chunk_chars]) i += paso devolver trozos def canonical_chunk_id(s: str) -> str: si s es Ninguno: devolver "" s = str(s).strip() s = s.strip("<>"'()[]{}") s = s.rstrip(".,;:") return s def inject_exec_summary_citations(exec_summary: str, citas: Lista[str], Allow_chunk_ids: Lista[str]) -> str: exec_summary = exec_summary o "" cset =[]para c en citas: c = canonical_chunk_id(c) si c y c en Allow_chunk_ids y c no en cset: cset.append(c) si len(cset) >= 2: romper si len(cset) < 2: para c en Allow_chunk_ids: si c no está en cset: cset.append(c) si len(cset) >= 2: romper si len(cset) >= 2: necesario = [c para c en cset si c no en exec_summary] si es necesario: exec_summary = exec_summary.strip() si exec_summary y no exec_summary.endswith("."): exec_summary += "." exec_summary += f" (cita: {cset[0]}) (cita: {cset[1]})" devolver exec_summary
Configuramos el entorno, cargamos de forma segura la clave API de OpenAI e inicializamos las utilidades principales de las que depende todo lo demás. Definimos hash, normalización de URL, limpieza de HTML y fragmentación para que todos los pasos posteriores funcionen en texto limpio y consistente. También agregamos ayudas deterministas para normalizar e inyectar citas, asegurando que las barreras de seguridad siempre se cumplan. Consulta los CÓDIGOS COMPLETOS aquí.
async def fetch_many(urls: List[str], timeout_s: float = 25.0, per_url_char_limit: int = 60000) -> Dict[str, str]: headers = {"User-Agent": "Mozilla/5.0 (AgenticAI/4.2)"} urls = [normalize_url(u) for u in urls] urls = [u para u en las URL si u.startswith("http")] urls = list(dict.fromkeys(urls)) out: Dict[str, str] = {} async con httpx.AsyncClient(timeout=timeout_s, follow_redirects=True, headers=headers) como cliente: async def _one(url: str): intente: r = await client.get(url) r.raise_for_status() out[url] = clean_html_to_text(r.text)[:per_url_char_limit] excepto excepción como e: out[url] = f"__FETCH_ERROR__ {type(e).__name__}: {e}" await asyncio.gather(*[_one(u) for u in urls]) return out def dedupe_texts(fuentes: Dict[str, str]) -> Dict[str, str]: visto = set() out = {} para url, txt en fuentes.items(): si no es instancia(txt, str) o txt.startswith("__FETCH_ERROR__"): continuar h = sha1(txt[:25000]) si h en visto: continuar visto.add(h) out[url] = txt devolver clase ChunkRecord(BaseModel): chunk_id: str url: str chunk_index: int texto: str clase RetrievalHit(BaseModel): chunk_id: str url: str chunk_index: int score_sparse: float = 0.0 score_dense: float = 0.0 score_fused: float = 0.0 texto: str clase EvidencePack(BaseModel): consulta: str hits: Lista[RetrievalHit]
Recuperamos de forma asincrónica múltiples fuentes web en paralelo y deduplicamos agresivamente el contenido para evitar evidencia redundante. Convertimos páginas sin formato en texto estructurado y definimos los modelos de datos centrales que representan fragmentos y visitas de recuperación. Nos aseguramos de que cada fragmento de texto sea rastreable hasta una fuente específica y un índice de fragmentos. Consulta los CÓDIGOS COMPLETOS aquí.
EPISODE_DB = "agentic_episode_memory.db" def episodio_db_init(): con = sqlite3.connect(EPISODE_DB) cur = con.cursor() cur.execute(""" CREAR TABLA SI NO EXISTE episodios ( id INTEGER PRIMARY KEY AUTOINCREMENT, ts INTEGER NOT NULL, question TEXT NOT NULL, urls_json TEXT NOT NULL, retrieval_queries_json TEXTO NO NULO, útiles_fuentes_json TEXTO NO NULO ) """) con.commit() con.close() def episodio_store(pregunta: cadena, URL: Lista[cadena], consultas_retrieval: Lista[cadena], fuentes_útiles: Lista[cadena]): con = sqlite3.connect(EPISODE_DB) cur = con.cursor() cur.execute( "INSERT EN episodios(ts, pregunta, urls_json, recuperación_queries_json, fuentes_útiles_json) VALORES(?,?,?,?,?)", (int(time.time()), pregunta, json.dumps(urls), json.dumps(recuperación_queries), json.dumps(fuentes_útiles)), ) con.commit() con.close() def episodio_recall(pregunta: str, top_k: int = 2) -> Lista[Dict[str, Any]]: con = sqlite3.connect(EPISODE_DB) cur = con.cursor() cur.execute("SELECT ts, question, urls_json, retrieval_queries_json, útil_sources_json FROM episodios ORDER BY ts DESC LIMIT 200") filas = cur.fetchall() con.close() q_tokens = set(re.findall(r"[A-Za-z]{3,}", (pregunta o "").lower())) puntuado =[]para ts, q2, u, rq, us en filas: t2 = set(re.findall(r"[A-Za-z]{3,}", (q2 o "").lower())) si no es t2: continuar puntuación = len(q_tokens & t2) / max(1, len(q_tokens)) si puntuación > 0: scoring.append((score, { "ts": ts, "question": q2, "urls": json.loads(u), "retrieval_queries": json.loads(rq), "useful_sources": json.loads(us), })) scoring.sort(key=lambda x: x[0], reverso=Verdadero) devuelve [x[1]para x en puntuado[:top_k]] episodio_db_init()
Introducimos la memoria episódica respaldada por SQLite para que el sistema pueda recordar lo que funcionó en ejecuciones anteriores. Almacenamos preguntas, estrategias de recuperación y fuentes útiles para guiar la planificación futura. También implementamos un recuerdo ligero basado en similitudes para sesgar el sistema hacia patrones históricamente efectivos. Consulta los CÓDIGOS COMPLETOS aquí.
clase HybridIndex: def __init__(self): self.records: Lista[ChunkRecord] =[]self.tfidf: Opcional[TfidfVectorizer] = Ninguno self.tfidf_mat = Ninguno self.emb_mat: Opcional[np.ndarray] = Ninguno def build_sparse(self): corpus = [r.text for r in self.records] if self.records else [""] self.tfidf = TfidfVectorizer(stop_words="english", ngram_range=(1, 2), max_features=80000) self.tfidf_mat = self.tfidf.fit_transform(corpus) def search_sparse(self, query: str, k: int) -> List[Tuple[int, float]]: si no self.records o self.tfidf es Ninguno o self.tfidf_mat es Ninguno: return[]qv = self.tfidf.transform([consulta]) sims = cosine_similarity(qv, self.tfidf_mat).flatten() top = np.argsort(-sims)[:k] return [(int(i), float(sims[i])) for i in top] def set_dense(self, mat: np.ndarray): self.emb_mat = mat.astype(np.float32) def search_dense(self, q_emb: np.ndarray, k: int) -> Lista[Tuple[int, float]]: si self.emb_mat es Ninguno o no self.records: return[]M = self.emb_mat q = q_emb.astype(np.float32).reshape(1, -1) M_norm = M / (np.linalg.norm(M, axis=1, keepdims=True) + 1e-9) q_norm = q / (np.linalg.norm(q) + 1e-9) sims = (M_norm @ q_norm.T).flatten() top = np.argsort(-sims)[:k] return [(int(i), float(sims[i])) for i in top] def rrf_fuse(rankings: List[List[int]], k: int = 60) -> Dict[int, float]: puntuaciones: Dict[int, float] = {} for r in rankings: for pos, idx in enumerar(r, inicio=1): puntuaciones[idx] = puntuaciones.get(idx, 0.0) + 1.0 / (k + pos) devolver puntuaciones HYBRID = HybridIndex() ALLOWED_URLS: Lista[str] =[]EMBED_MODEL = "text-embedding-3-small" async def embed_batch(textos: Lista[str]) -> np.ndarray: resp = await oa.embeddings.create(model=EMBED_MODEL, input=texts, encoding_format="float") vecs = [np.array(item.embedding, dtype=np.float32) para elemento en resp.data] devuelve np.vstack(vecs) si vecs else np.zeros((0, 0), dtype=np.float32) async def embed_texts(texts: List[str], batch_size: int = 96, max_concurrency: int = 3) -> np.ndarray: sem = asyncio.Semaphore(max_concurrency) tapetes: Lista[Tupla[int, np.ndarray]] =[]async def _one(inicio: int, lote: Lista[str]): async con sem: m = espera embed_batch(batch) mats.append((inicio, m)) tareas =[]para comenzar en rango (0, len (textos), tamaño de lote): lote = [t[:7000] para t en textos [inicio: inicio + tamaño de lote]] tareas.append(_one(inicio, lote)) await asyncio.gather(*tasks) mats.sort(key=lambda x: x[0]) emb = np.vstack([m para _, m en tapetes]) si tapetes else np.zeros((len(textos), 0), dtype=np.float32) si emb.shape[0]!= len(texts): rise RuntimeError(f"Las filas incrustadas no coinciden: tengo {emb.shape[0]} esperado {len(textos)}") return emb async def embed_query(query: str) -> np.ndarray: m = await embed_batch([query[:7000]]) return m[0]si m.forma[0]else np.zeros((0,), dtype=np.float32) async def build_index(urls: List[str], max_chunks_per_url: int = 60): global ALLOWED_URLS fetched = await fetch_many(urls) fetched = dedupe_texts(fetched) registros: List[ChunkRecord] =[]permitido: Lista[cadena] =[]para url, txt en fetched.items(): si no es isinstance(txt, str) o txt.startswith("__FETCH_ERROR__"): continúa permitido.append(url) fragmentos = chunk_text(txt)[:max_chunks_per_url] para i, ch en enumerate(fragmentos): cid = f"{sha1(url)}:{i}" records.append(ChunkRecord(chunk_id=cid, url=url, chunk_index=i, text=ch)) si no registros: err_view = {normalize_url(u): fetched.get(normalize_url(u), "") for u en urls} rise RuntimeError("No se han obtenido fuentes exitosamente.n" + json.dumps(err_view, sangría=2)[:4000]) ALLOWED_URLS = permitido HYBRID.records = registros HYBRID.build_sparse() textos = [r.text para r en HYBRID.records] emb = espera embed_texts(textos, tamaño de lote=96, max_concurrency=3) HYBRID.set_dense(emb)
Creamos un índice de recuperación híbrido que combina una búsqueda TF-IDF escasa con incrustaciones densas de OpenAI. Permitimos la fusión de rangos recíproca, de modo que las señales escasas y densas se complementen entre sí en lugar de competir. Construimos el índice una vez por ejecución y lo reutilizamos en todas las consultas de recuperación para mayor eficiencia. Consulta los CÓDIGOS COMPLETOS aquí.
def build_evidence_pack(consulta: str, sparse: Lista[Tuple[int,float]], denso: Lista[Tuple[int,float]], k: int = 10) -> EvidencePack: sparse_rank = [i para i,_ en disperso] denso_rank = [i para i,_ en denso] sparse_scores = {i:s para i,s en disperso} denso_scores = {i:s para i,s en denso} fusionado = rrf_fuse([sparse_rank, denso_rank], k=60) si denso_rank else rrf_fuse([sparse_rank], k=60) top = sorted(fused.keys(), key=lambda i: fusionado[i], reverso=True)[:k] visitas: Lista[RetrievalHit] =[]para idx en la parte superior: r = HYBRID.records[idx] hits.append(RetrievalHit( chunk_id=r.chunk_id, url=r.url, chunk_index=r.chunk_index, score_sparse=float(sparse_scores.get(idx, 0.0)), score_dense=float(dense_scores.get(idx, 0.0)), score_fused=float(fused.get(idx, 0.0)), text=r.text )) return EvidencePack(query=consulta, hits=hits) async def together_evidence(consultas: Lista[str], per_query_k: int = 10, sparse_k: int = 60, densa_k: int = 60): evidencia: Lista[EvidencePack] =[]cuenta_fuentes_útiles: Dict[cadena, int] = {} all_chunk_ids: Lista[cadena] =[]para q en consultas: sparse = HYBRID.search_sparse(q, k=sparse_k) q_emb = await embed_query(q) denso = HYBRID.search_dense(q_emb, k=dense_k) pack = build_evidence_pack(q, sparse, denso, k=per_query_k) evidencia.append(pack) para h en pack.hits[:6]: cuenta_fuentes_útiles[h.url] = cuenta_fuentes_útiles.get(h.url, 0) + 1 para h en pack.hits: all_chunk_ids.append(h.chunk_id) fuentes_útiles = ordenado(fuentes_útiles_count.keys(), clave=lambda u: cuenta_fuentes_útiles[u], reverso=True) all_chunk_ids = ordenado (lista (dict.fromkeys (all_chunk_ids))) evidencia de devolución, fuentes_útiles [: 8], clase all_chunk_ids Plan (BaseModel): objetivo: str subtareas: Lista [str] consultas_de recuperación: Lista [str] comprobaciones_de aceptación: Lista [str] clase UltraAnswer (BaseModel): título: str resumen_ejecutivo: arquitectura str: Lista [str] estrategia_de recuperación: Lista[str] agent_graph: Lista[str] notas_de_implementación: Lista[str] riesgos_y_limites: Lista[str] citas: Lista[str] fuentes: Lista[str] def normalize_answer(ans: UltraAnswer, Allow_chunk_ids: Lista[str]) -> UltraAnswer: datos = ans.model_dump() datos["citaciones"] = [canonical_chunk_id(x) para x en (data.get("citas") o[])] datos["citaciones"] = [x para x en datos["citaciones"] si x en permitido_chunk_ids] datos["executive_summary"] = inject_exec_summary_citations(data.get("executive_summary",""), datos["citations"], permitido_chunk_ids) return UltraAnswer(**data) def validar_ultra(ans: UltraAnswer, permitido_chunk_ids: Lista[str]) -> Ninguno: extras = [u para u en ans.sources si no está en ALLOWED_URLS] si extras: elevar ValueError(f"Fuentes no permitidas en la salida: {extras}") cset = set(ans.citations o[]) desaparecido = [cid para cid en cset si cid no está en conjunto(allowed_chunk_ids)] si falta: elevar ValueError(f"Las citas hacen referencia a fragmentos_ids desconocidos (no recuperados): {missing}") si len(cset) < 6: elevar ValueError("Necesita al menos 6 citas distintas de identificadores de fragmentos en modo ultra.") es_text = ans.executive_summary o "" es_count = sum(1 for cid in cset if cid in es_text) if es_count < 2: rise ValueError("El resumen ejecutivo debe incluir al menos 2 citas de chunk_id textualmente.") PLANNER = Agent( name="Planner", model="gpt-4o-mini", instrucciones=( "Devolver un esquema de plan técnico.n" "Hacer entre 10 y 16 consultas_de recuperación.n" "La aceptación debe incluir: al menos 6 citas y exec_summary contiene al menos 2 citas palabra por palabra." ), tipo de salida=Plan, ) SINTETIZADOR = Agente( nombre="Sintetizador", modelo="gpt-4o-mini", instrucciones=( "Devolver el esquema de UltraAnswer.n" "Restricciones estrictas:n" "- Executive_summary DEBE incluir al menos DOS citas palabra por palabra como: (cita: ).n" "- las citas deben elegirse ÚNICAMENTE de la lista ALLOWED_CHUNK_IDS.n" "- la lista de citas debe incluir al menos 6 fragment_ids únicos.n" "- las fuentes deben ser un subconjunto de URL permitidas.n" ), output_type=UltraAnswer, ) FIXER = Agent( name="Fixer", model="gpt-4o-mini", instrucciones=( "Reparar para satisfacer las barreras de seguridad.n" "Asegúrese de que Executive_summary incluya al menos DOS citas textualmente.n" "Elija citas SÓLO de la lista ALLOWED_CHUNK_IDS.n" "Devuelva el esquema de UltraAnswer." ), output_type=UltraAnswer, ) session = SQLiteSession("ultra_agentic_user", "ultra_agentic_session.db")
Recopilamos evidencia ejecutando múltiples consultas específicas, fusionando resultados dispersos y densos y reuniendo paquetes de evidencia con puntuaciones y procedencia. Definimos esquemas estrictos para los planes y las respuestas finales, luego normalizamos y validamos las citas con los ID de fragmentos recuperados. Aplicamos barreras estrictas para que cada respuesta permanezca fundamentada y sea auditable. Consulta los CÓDIGOS COMPLETOS aquí.
async def run_ultra_agentic(pregunta: str, urls: List[str], max_repairs: int = 2) -> UltraAnswer: await build_index(urls) retirada_hint = json.dumps(episode_recall(pregunta, top_k=2), indent=2)[:2000] plan_res = await Runner.run( PLANNER, f"Pregunta:n{pregunta}nnURL permitidas:n{json.dumps(ALLOWED_URLS, indent=2)}nnRecall:n{recall_hint}n", session=session ) plan: Plan = plan_res.final_output consultas = (plan.retrieval_queries o[])[:16] evidencia_packs, fuentes_útiles, permitido_chunk_ids = aguardar recopilación_evidencia(consultas) evidencia_json = json.dumps([p.model_dump() para p en evidencia_packs], sangría=2)[:16000] permitido_chunk_ids_json = json.dumps(allowed_chunk_ids[:200], sangría=2) draft_res = await Runner.run( SYNTHESIZER, f"Pregunta:n{pregunta}nnURL permitidas:n{json.dumps(ALLOWED_URLS, indent=2)}nn" f"ALLOWED_CHUNK_IDS:n{allowed_chunk_ids_json}nn" f"Paquetes de evidencia:n{evidence_json}nn" "Devolver UltraAnswer.", sesión=sesión) borrador = normalize_answer(borrador_res.final_output, permitido_chunk_ids) last_err = Ninguno para i en el rango(max_repairs + 1): intente: validar_ultra(borrador, permitido_chunk_ids) episodio_store(pregunta, ALLOWED_URLS, plan.retrieval_queries, fuentes_útiles) devolver borrador excepto Excepción como e: last_err = str(e) if i >= max_repairs: borrador = normalize_answer(borrador, permitido_chunk_ids) validar_ultra(borrador, permitido_chunk_ids) devolver borrador fixer_res = await Runner.run( FIXER, f"Pregunta:n{pregunta}nnURL permitidas:n{json.dumps(ALLOWED_URLS, indent=2)}nn" f"ALLOWED_CHUNK_IDS:n{allowed_chunk_ids_json}nn" f"Error de barandilla:n{last_err}nn" f"Borrador:n{json.dumps(draft.model_dump(), indent=2)[:12000]}nn" f"Paquetes de evidencia:n{evidence_json}nn" "Devolver UltraAnswer corregido que pasa las barreras de seguridad.", session=session ) draft = normalize_answer(fixer_res.final_output, permitido_chunk_ids) rise RuntimeError(f"Fallo inesperado: {last_err}") question = ( "Diseñe un flujo de trabajo de IA agente avanzado pero de producción eficiente en Python con recuperación híbrida, " "citas de procedencia primero, ciclos de crítica y reparación, y memoria episódica. " "Explique por qué es importante cada capa, modos de falla y evaluación." ) urls = [ "https://openai.github.io/openai-agents-python/", "https://openai.github.io/openai-agents-python/agents/", "https://openai.github.io/openai-agents-python/running_agents/", "https://github.com/openai/openai-agents-python", ] ans = await run_ultra_agentic(question, urls, max_repairs=2) print("nTITLE:n", ans.title) print("nRESUMEN EJECUTIVO:n", ans.executive_summary) print("nARQUITECTURA:") para x en ans.architecture: print("-", x) print("nESTRATEGIA DE RECUPERACIÓN:") para x en ans.retrieval_strategy: print("-", x) print("nGRÁFICO DEL AGENTE:") para x en ans.agent_graph: print("-", x) print("nNOTAS DE IMPLEMENTACIÓN:") para x en ans.implementation_notes: print("-", x) print("nRIESGOS & LÍMITES:") para x en ans.risks_and_limits: print("-", x) print("nCITACIONES (chunk_ids):") para c en ans.citations: print("-", c) print("nFUENTES:") para s en ans.sources: print("-", s)
Orquestamos el ciclo agente completo encadenando la planificación, la síntesis, la validación y la reparación en una canalización segura asíncrona. Reintentamos y reparamos automáticamente los resultados hasta que superan todas las restricciones sin intervención humana. Terminamos ejecutando un ejemplo completo e imprimiendo una respuesta agente completamente fundamentada y lista para producción.
En conclusión, desarrollamos un canal de agencia integral y resistente a los modos de falla comunes: formas de incrustación inestables, deriva de citas y falta de conexión a tierra en los resúmenes ejecutivos. Validamos los resultados con fuentes incluidas en la lista permitida, recuperamos ID de fragmentos, normalizamos automáticamente las citas e inyectamos citas deterministas cuando fue necesario para garantizar el cumplimiento sin sacrificar la corrección. Al combinar recuperación híbrida, ciclos de crítica y reparación y memoria episódica, creamos una base reutilizable que podemos ampliar con evaluaciones más sólidas (puntuación de cobertura de reclamo a evidencia, equipos rojos adversarios y pruebas de regresión) para fortalecer continuamente el sistema a medida que escala a nuevos dominios y corpus más grandes.
Consulta los CÓDIGOS COMPLETOS aquí. Además, no dude en seguirnos en Twitter y no olvide unirse a nuestro SubReddit de más de 100.000 ML y suscribirse a nuestro boletín. ¡Esperar! estas en telegrama? Ahora también puedes unirte a nosotros en Telegram.