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
5 changes: 5 additions & 0 deletions changelog.d/features/orchestration-agents-ws.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
- **feat(dashboard):** the `/dashboard/orchestration` snapshot hook now subscribes to the
`agents` WebSocket channel (`agent.task.updated`) instead of `requests` as its refetch
trigger, and relaxes its background poll from 5s to 30s while that WS connection is up —
falling back to the tighter 5s cadence, reprogrammed live on any connect/disconnect
transition, whenever the socket is down.
Original file line number Diff line number Diff line change
@@ -1,5 +1,9 @@
"use client";
/** Polls the 3 agent sources (allSettled), listens to the `requests` WS channel as a refetch trigger. */
/**
* Polls the 3 agent sources (allSettled), listens to the `agents` WS channel as a refetch
* trigger, and relaxes the poll interval from 5s to 30s while that WS connection is up (the
* channel event still forces an immediate debounced refetch either way).
*/
import { useCallback, useEffect, useMemo, useRef, useState } from "react";
import { useLiveDashboard } from "@/hooks/useLiveDashboard";
import type { CloudAgentTask } from "@/lib/cloudAgent/types";
Expand All @@ -12,6 +16,7 @@ import { mergeSnapshot } from "../model/mergeSnapshot";
import type { OrchSnapshot, SourceStatus } from "../model/orchestrationTypes";

export const POLL_MS = 5_000;
export const POLL_MS_WS_CONNECTED = 30_000;
export const WS_REFETCH_DEBOUNCE_MS = 1_000;

interface Raw {
Expand Down Expand Up @@ -135,9 +140,7 @@ export function useOrchestrationSnapshot() {

pollRef.current = () => void poll();
void poll();
const id = setInterval(() => void poll(), POLL_MS);
return () => {
clearInterval(id);
controller.abort();
if (debounceRef.current) {
clearTimeout(debounceRef.current);
Expand All @@ -150,10 +153,10 @@ export function useOrchestrationSnapshot() {
pollRef.current();
}, []);

useLiveDashboard({
channels: ["requests"],
const { connection } = useLiveDashboard({
channels: ["agents"],
onEvent: (payload) => {
if (payload.channel !== "requests") return;
if (payload.channel !== "agents") return;
if (debounceRef.current) return; // debounce burst → one refetch
debounceRef.current = setTimeout(() => {
debounceRef.current = null;
Expand All @@ -162,6 +165,17 @@ export function useOrchestrationSnapshot() {
},
});

// Adaptive poll interval, separated from the mount effect above: a connected `agents` WS
// already pushes refetches on change, so the background poll only needs to be a slow safety
// net (30s) — it falls back to the tighter 5s cadence while the WS is down. Declared AFTER
// the mount effect so `pollRef.current` is already populated (its initial value is a safe
// no-op) by the time this effect's first tick can fire.
const wsConnected = connection.isConnected;
useEffect(() => {
const id = setInterval(() => pollRef.current(), wsConnected ? POLL_MS_WS_CONNECTED : POLL_MS);
return () => clearInterval(id);
}, [wsConnected]);

const snapshot: OrchSnapshot = useMemo(
() =>
mergeSnapshot(
Expand Down
17 changes: 17 additions & 0 deletions src/lib/a2a/taskManager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,20 @@

import { randomUUID } from "crypto";

import { emit } from "@/lib/events/eventBus";

/**
* Publish an `agent.task.updated` transition for the orchestration canvas (Fase 2, Task B2).
* Best-effort: a listener throwing must never break the task write path that triggered it.
*/
function emitAgentTaskUpdated(source: "cloud-agent" | "a2a", taskId: string, state: string): void {
try {
emit("agent.task.updated", { source, taskId, state, timestamp: Date.now() });
} catch {
/* listeners never derail the write path */
}
}

// ============ Types ============

export type TaskState = "submitted" | "working" | "completed" | "failed" | "cancelled";
Expand Down Expand Up @@ -114,6 +128,7 @@ export class A2ATaskManager {
...(owner !== undefined ? { owner } : {}),
};
this.tasks.set(task.id, task);
emitAgentTaskUpdated("a2a", task.id, "submitted");
return task;
}

Expand Down Expand Up @@ -158,6 +173,7 @@ export class A2ATaskManager {
task.events.push({ timestamp: now, state, message });
if (artifacts) task.artifacts.push(...artifacts);

emitAgentTaskUpdated("a2a", taskId, state);
return task;
}

Expand Down Expand Up @@ -243,6 +259,7 @@ export class A2ATaskManager {
task.state = "failed";
task.updatedAt = now.toISOString();
task.events.push({ timestamp: now.toISOString(), state: "failed", message: "TTL expired" });
emitAgentTaskUpdated("a2a", id, "failed");
}
// Remove terminal tasks older than 2x TTL
if (
Expand Down
15 changes: 15 additions & 0 deletions src/lib/cloudAgent/db.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,17 @@
import { getDbInstance } from "@/lib/db/core.ts";
import { emit } from "@/lib/events/eventBus";

/**
* Publish an `agent.task.updated` transition for the orchestration canvas (Fase 2, Task B2).
* Best-effort: a listener throwing must never break the DB write path that triggered it.
*/
function emitAgentTaskUpdated(source: "cloud-agent" | "a2a", taskId: string, state: string): void {
try {
emit("agent.task.updated", { source, taskId, state, timestamp: Date.now() });
} catch {
/* listeners never derail the write path */
}
}

export interface CloudAgentTaskRow {
id: string;
Expand Down Expand Up @@ -66,6 +79,7 @@ export function insertCloudAgentTask(task: CloudAgentTaskRow): void {
)
`
).run(task);
emitAgentTaskUpdated("cloud-agent", task.id, task.status);
}

// Whitelist of allowed columns for update operations
Expand Down Expand Up @@ -107,6 +121,7 @@ export function updateCloudAgentTask(
WHERE id = @id
`
).run({ id, ...validUpdates });
emitAgentTaskUpdated("cloud-agent", id, (validUpdates.status as string) ?? "updated");
}

export function getCloudAgentTaskById(id: string): CloudAgentTaskRow | null {
Expand Down
40 changes: 40 additions & 0 deletions src/lib/conductor/hubProxy.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,8 @@

import { z } from "zod";

import { emit } from "@/lib/events/eventBus";

// ============ Whitelisted client-facing shapes ============

export interface FleetRunner {
Expand Down Expand Up @@ -112,6 +114,43 @@ function toFleetTask(t: z.infer<typeof hubTaskSchema>): FleetTask {
};
}

// ============ Fleet task mirror (Orchestration Canvas Fase 2, Task B3) ============
//
// Module-level cache of the last known status per fleet task, so `getFleetSnapshot` can
// diff-on-fetch and mirror Conductor task transitions into the `agents` WS channel without a
// dedicated poller — it piggybacks on the dashboard's existing poll. `null` means "no snapshot
// observed yet" (first-ever call): that call only seeds the cache, it never emits, since the
// dashboard already fetches the full snapshot on its initial poll. An offline snapshot never
// touches this cache (see call site below), so a hub flap does not cause a re-seed burst once
// the hub comes back — only the real delta since the last successful snapshot is emitted.
let lastFleetTaskStates: Map<string, string> | null = null;

function emitFleetTransitions(tasks: FleetTask[]): void {
const next = new Map(tasks.map((t) => [t.id, t.status]));
if (lastFleetTaskStates) {
for (const [id, status] of next) {
if (lastFleetTaskStates.get(id) !== status) {
try {
emit("agent.task.updated", {
source: "conductor",
taskId: id,
state: status,
timestamp: Date.now(),
});
} catch {
/* best-effort */
}
}
}
}
lastFleetTaskStates = next;
}

/** Test-only seam: resets the fleet task mirror cache so tests are order-independent. */
export function __resetFleetMirrorForTests(): void {
lastFleetTaskStates = null;
}

/** Fleet snapshot for the dashboard panel. Degraded ({offline: true}) on any failure. */
export async function getFleetSnapshot(opts: HubProxyOptions = {}): Promise<FleetSnapshot> {
try {
Expand All @@ -128,6 +167,7 @@ export async function getFleetSnapshot(opts: HubProxyOptions = {}): Promise<Flee
draining: r.draining === true,
}));
const tasks = z.array(hubTaskSchema).parse(rawTasks).map(toFleetTask);
emitFleetTransitions(tasks);
return { offline: false, runners, tasks };
} catch {
return { offline: true, runners: [], tasks: [] };
Expand Down
19 changes: 17 additions & 2 deletions src/lib/events/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,8 @@ export type DashboardEventName =
| "combo.target.succeeded"
| "credential.health.changed"
| "compression.completed"
| "compression.step";
| "compression.step"
| "agent.task.updated";

// ── Event Payloads ────────────────────────────────────────────────────────

Expand Down Expand Up @@ -138,6 +139,18 @@ export interface CompressionStepPayload {
timestamp: number;
}

/**
* Real-time transition of an orchestrated task tracked by the orchestration dashboard
* (Fase 2). `source` distinguishes which subsystem raised the transition so the canvas can
* route it to the right node type.
*/
export interface AgentTaskUpdatedPayload {
source: "cloud-agent" | "a2a" | "conductor";
taskId: string;
state: string;
timestamp: number;
}

// ── Event Map ─────────────────────────────────────────────────────────────

export interface DashboardEventMap {
Expand All @@ -151,6 +164,7 @@ export interface DashboardEventMap {
"credential.health.changed": CredentialHealthChangedPayload;
"compression.completed": CompressionCompletedPayload;
"compression.step": CompressionStepPayload;
"agent.task.updated": AgentTaskUpdatedPayload;
}

// ── Event Bus Listener ────────────────────────────────────────────────────
Expand All @@ -162,14 +176,15 @@ export type DashboardEventListener<E extends DashboardEventName> = (
// ── Channel Definitions ───────────────────────────────────────────────────

/** Available subscription channels */
export type DashboardChannel = "requests" | "combo" | "credentials" | "compression";
export type DashboardChannel = "requests" | "combo" | "credentials" | "compression" | "agents";

/** Map channels to their events */
export const CHANNEL_EVENTS: Record<DashboardChannel, DashboardEventName[]> = {
requests: ["request.started", "request.streaming", "request.completed", "request.failed"],
combo: ["combo.target.attempt", "combo.target.failed", "combo.target.succeeded"],
credentials: ["credential.health.changed"],
compression: ["compression.completed", "compression.step"],
agents: ["agent.task.updated"],
};

/** Get channel for an event */
Expand Down
6 changes: 3 additions & 3 deletions src/server/ws/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@

export interface WsSubscribeMessage {
type: "subscribe";
channels: Array<"requests" | "combo" | "credentials" | "compression">;
channels: Array<"requests" | "combo" | "credentials" | "compression" | "agents">;
}

export interface WsPingMessage {
Expand All @@ -21,7 +21,7 @@ export type WsClientMessage = WsSubscribeMessage | WsPingMessage;

export interface WsEventMessage {
type: "event";
channel: "requests" | "combo" | "credentials" | "compression";
channel: "requests" | "combo" | "credentials" | "compression" | "agents";
event: string;
data: unknown;
}
Expand All @@ -35,7 +35,7 @@ export interface WsWelcomeMessage {
version: string;
sessionId: string;
serverTime: number;
channels: Array<"requests" | "combo" | "credentials" | "compression">;
channels: Array<"requests" | "combo" | "credentials" | "compression" | "agents">;
/** Number of buffered events since last reconnect */
backlog: number;
}
Expand Down
Loading