File size: 2,451 Bytes
03e5649
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
"""
speculative_engine.py — S-SPECULATIVE: Multi-Relay Parallelism.
Lancia tentativi multipli in parallelo con diverse strategie e seleziona il primo vincente.
Ottimizza la latenza e la resilienza contro i blocchi logici.
"""
import asyncio
import logging
from typing import List, Dict, Any, Callable, Awaitable

_logger = logging.getLogger("agente_ai.speculative")

class SpeculativeEngine:
    def __init__(self, verifier_fn: Callable[[str], Awaitable[bool]]):
        self.verifier_fn = verifier_fn

    async def run_speculative(self, tasks: List[Awaitable[str]], timeout: float = 60.0) -> str:
        """
        Esegue più task in parallelo. Il primo che restituisce una risposta 
        che passa la verifica viene accettato. Gli altri vengono cancellati.
        """
        if not tasks:
            return ""

        # Creiamo i task asyncio
        pending = [asyncio.create_task(t) for t in tasks]
        
        try:
            while pending:
                # Aspettiamo che il primo task finisca
                done, pending = await asyncio.wait(
                    pending, 
                    return_when=asyncio.FIRST_COMPLETED,
                    timeout=timeout
                )
                
                if not done: # Timeout
                    break

                for task in done:
                    try:
                        result = await task
                        # Validazione speculativa
                        if await self.verifier_fn(result):
                            _logger.info("[Speculative] Soluzione vincente trovata! Cancellazione altri task.")
                            # Cancella i task ancora in corso
                            for p in pending:
                                p.cancel()
                            return result
                        else:
                            _logger.debug("[Speculative] Task completato ma non ha superato la verifica.")
                    except Exception as e:
                        _logger.error(f"[Speculative] Errore in un task parallelo: {e}")

            return "Errore: Nessun task speculativo ha prodotto una soluzione valida."
        
        finally:
            # Pulizia finale
            for p in pending:
                p.cancel()

# Esempio di utilizzo nel loop:
# engine = SpeculativeEngine(verifier_fn=state.verifier.verify)
# winner = await engine.run_speculative([attempt1, attempt2, attempt3])