Skip to content
Open
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
10 changes: 5 additions & 5 deletions backend/src/Agent/TurnSignal.hs
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ module Agent.TurnSignal (
signalTurnLog,
) where

import Control.Concurrent.STM (TChan, TVar, atomically, modifyTVar', newTChan, newTVarIO, readTVar, writeTChan, writeTVar)
import Control.Concurrent.STM (TChan, TVar, atomically, dupTChan, modifyTVar', newBroadcastTChan, newTVarIO, readTVar, writeTChan, writeTVar)
import Data.Map.Strict (Map)
import qualified Data.Map.Strict as Map
import System.FilePath (normalise)
Expand All @@ -21,19 +21,19 @@ import System.IO.Unsafe (unsafePerformIO)
signals :: TVar (Map FilePath (TChan ()))
signals = unsafePerformIO $ newTVarIO Map.empty

{- | Get (or create) the wakeup channel for a turn log. Concurrent
streams for the same turn share one channel.
{- | Get a private wakeup channel for a turn log.
-}
registerTurnSignal :: FilePath -> IO (TChan ())
registerTurnSignal path = atomically $ do
let key = normalise path
m <- readTVar signals
case Map.lookup key m of
broadcast <- case Map.lookup key m of
Just ch -> pure ch
Nothing -> do
ch <- newTChan
ch <- newBroadcastTChan
writeTVar signals (Map.insert key ch m)
pure ch
dupTChan broadcast

{- | Stop tracking a turn log. Idempotent; safe to call when the turn
finishes even if no stream ever connected.
Expand Down
64 changes: 43 additions & 21 deletions frontend/ffi.js
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,10 @@ let stepStatusTargetKey = null;
let stepStatusApplicationError = false;
let agentTurnSource = null;
let agentTurnTargetKey = null;
const AGENT_TURN_INITIAL_RETRY_DELAY = 1000;
const AGENT_TURN_MAX_RETRY_DELAY = 30000;
let agentTurnDeliveredLength = 0;
let agentTurnRetryDelay = AGENT_TURN_INITIAL_RETRY_DELAY;

function closeStepStatusStream() {
if (stepStatusSource) {
Expand All @@ -78,6 +82,8 @@ function closeAgentTurnStream() {
agentTurnSource = null;
}
agentTurnTargetKey = null;
agentTurnDeliveredLength = 0;
agentTurnRetryDelay = AGENT_TURN_INITIAL_RETRY_DELAY;
}

function toggleTheme() {
Expand Down Expand Up @@ -185,31 +191,27 @@ export function connectPorts(app) {
};
}

function openAgentTurnStream({ turnId }) {
const url = `/backend/agent/turn/${encodeURIComponent(turnId)}/stream`;

if (
agentTurnSource &&
agentTurnTargetKey === url &&
agentTurnSource.readyState !== EventSource.CLOSED
) {
return;
}

closeAgentTurnStream();
function connectAgentTurnStream(turnId) {
const source = new EventSource(
`/backend/agent/turn/${encodeURIComponent(turnId)}/stream`,
);
let replayLength = agentTurnDeliveredLength;
agentTurnSource = source;

agentTurnSource = new EventSource(url);
agentTurnTargetKey = url;

agentTurnSource.addEventListener("chunk", (event) => {
source.addEventListener("chunk", (event) => {
try {
emitAgentTurn("chunk", JSON.parse(event.data));
const data = JSON.parse(event.data);
data.chunk = data.chunk.slice(replayLength);
replayLength = 0;
agentTurnDeliveredLength += data.chunk.length;
agentTurnRetryDelay = AGENT_TURN_INITIAL_RETRY_DELAY;
if (data.chunk) emitAgentTurn("chunk", data);
} catch (err) {
emitAgentTurn("error", `Failed to parse agent log chunk: ${String(err)}`);
}
});

agentTurnSource.addEventListener("done", (event) => {
source.addEventListener("done", (event) => {
try {
emitAgentTurn("done", JSON.parse(event.data));
} catch {
Expand All @@ -218,17 +220,37 @@ export function connectPorts(app) {
closeAgentTurnStream();
});

agentTurnSource.addEventListener("heartbeat", (event) => {
source.addEventListener("heartbeat", (event) => {
agentTurnRetryDelay = AGENT_TURN_INITIAL_RETRY_DELAY;
try {
emitAgentTurn("heartbeat", JSON.parse(event.data));
} catch {}
});

agentTurnSource.onerror = () => {
emitAgentTurn("error", "Agent turn stream connection issue");
source.onerror = () => {
source.close();
if (source !== agentTurnSource || agentTurnTargetKey !== turnId) return;

agentTurnSource = null;
const delay = agentTurnRetryDelay;
if (delay === AGENT_TURN_INITIAL_RETRY_DELAY) {
emitAgentTurn("error", "Agent turn stream connection issue");
}
agentTurnRetryDelay = Math.min(delay * 2, AGENT_TURN_MAX_RETRY_DELAY);
setTimeout(() => {
if (agentTurnTargetKey === turnId && !agentTurnSource) {
connectAgentTurnStream(turnId);
}
}, delay);
};
}

function openAgentTurnStream({ turnId }) {
closeAgentTurnStream();
agentTurnTargetKey = turnId;
connectAgentTurnStream(turnId);
}

function storeLastChat(sessionId) {
localStorage.setItem("agent:lastChat", sessionId);
}
Expand Down