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
14 changes: 13 additions & 1 deletion apps/harness/lib/harness/metrics.ex
Original file line number Diff line number Diff line change
Expand Up @@ -11,9 +11,12 @@ defmodule Harness.Metrics do
Collect a full metrics snapshot.
"""
def collect do
sessions = session_metrics()

%{
beam: beam_metrics(),
sessions: session_metrics(),
sessions: sessions,
lifecycle: lifecycle_metrics(sessions),
snapshot_server: snapshot_server_metrics(),
timestamp: System.system_time(:millisecond)
}
Expand Down Expand Up @@ -85,6 +88,15 @@ defmodule Harness.Metrics do
end)
end

defp lifecycle_metrics(sessions) do
%{
active_sessions: length(sessions),
sessions_by_provider: Enum.frequencies_by(sessions, & &1.provider),
sessions_with_backlog: Enum.count(sessions, &(&1.message_queue_len > 0)),
total_message_queue_len: Enum.reduce(sessions, 0, &(&1.message_queue_len + &2))
}
end

defp process_metrics(pid) do
case :erlang.process_info(pid, [
:memory,
Expand Down
27 changes: 27 additions & 0 deletions apps/harness/lib/harness/provider_session.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
defmodule Harness.ProviderSession do
@moduledoc """
Behaviour implemented by provider-backed session processes.

`Harness.SessionManager` dispatches to this surface regardless of whether the
provider uses a persistent stdio process, a turn-scoped CLI, or an HTTP/SSE
bridge.
"""

@type request_id :: term()
@type decision :: term()
@type answers :: term()
@type params :: map()
@type thread_id :: String.t()
@type turn_id :: term()
@type snapshot :: term()

@callback start_link(keyword() | map()) :: GenServer.on_start()
@callback wait_for_ready(GenServer.server(), timeout()) :: :ok | {:error, term()}
@callback send_turn(GenServer.server(), params()) :: term()
@callback interrupt_turn(GenServer.server(), thread_id(), turn_id() | nil) :: term()
@callback respond_to_approval(GenServer.server(), request_id(), decision()) :: term()
@callback respond_to_user_input(GenServer.server(), request_id(), answers()) :: term()
@callback read_thread(GenServer.server(), thread_id()) :: {:ok, snapshot()} | {:error, term()}
@callback rollback_thread(GenServer.server(), thread_id(), non_neg_integer()) ::
{:ok, snapshot()} | {:error, term()}
end
1 change: 1 addition & 0 deletions apps/harness/lib/harness/providers/claude_session.ex
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ defmodule Harness.Providers.ClaudeSession do
parse its stdout JSON stream ourselves — same wire format, no SDK wrapper needed.
"""
use GenServer, restart: :temporary
@behaviour Harness.ProviderSession

alias Harness.Event

Expand Down
1 change: 1 addition & 0 deletions apps/harness/lib/harness/providers/codex_session.ex
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ defmodule Harness.Providers.CodexSession do
- Collab child conversation suppression
"""
use GenServer, restart: :temporary
@behaviour Harness.ProviderSession

alias Harness.JsonRpc
alias Harness.Event
Expand Down
1 change: 1 addition & 0 deletions apps/harness/lib/harness/providers/cursor_session.ex
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ defmodule Harness.Providers.CursorSession do
- Multi-turn via --resume flag
"""
use GenServer, restart: :temporary
@behaviour Harness.ProviderSession

alias Harness.Event

Expand Down
1 change: 1 addition & 0 deletions apps/harness/lib/harness/providers/mock_session.ex
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ defmodule Harness.Providers.MockSession do
- delayMs: delay between deltas in ms (default: 10)
"""
use GenServer, restart: :temporary
@behaviour Harness.ProviderSession

alias Harness.JsonRpc
alias Harness.Event
Expand Down
18 changes: 14 additions & 4 deletions apps/harness/lib/harness/providers/opencode_session.ex
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@ defmodule Harness.Providers.OpenCodeSession do
- Emits `collab_agent_spawn_begin` with parent→child linkage
"""
use GenServer, restart: :temporary
@behaviour Harness.ProviderSession

alias Harness.Event

Expand Down Expand Up @@ -564,7 +565,7 @@ defmodule Harness.Providers.OpenCodeSession do
{:line, 65_536},
{:cd, to_charlist(cwd)},
{:args, Enum.map(args, &to_charlist/1)},
{:env, build_env()}
{:env, build_env(state)}
]
)

Expand All @@ -574,9 +575,18 @@ defmodule Harness.Providers.OpenCodeSession do
end
end

defp build_env do
System.get_env()
|> Enum.map(fn {k, v} -> {to_charlist(k), to_charlist(v)} end)
defp build_env(state) do
base =
System.get_env()
|> Enum.map(fn {k, v} -> {to_charlist(k), to_charlist(v)} end)

opencode_config_path = get_in(state.params, ["providerOptions", "opencode", "configPath"])

if is_binary(opencode_config_path) and opencode_config_path != "" do
[{~c"OPENCODE_CONFIG", to_charlist(opencode_config_path)} | base]
else
base
end
end

defp wait_for_server(_base_url, 0), do: {:error, :timeout}
Expand Down
17 changes: 17 additions & 0 deletions apps/harness/test/harness/providers/codex_session_test.exs
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
defmodule Harness.Providers.CodexSessionTest do
use ExUnit.Case, async: true

alias Harness.Providers.CodexSession

test "exports the provider session callbacks used by SessionManager" do
assert function_exported?(CodexSession, :start_link, 1)
assert function_exported?(CodexSession, :wait_for_ready, 1)
assert function_exported?(CodexSession, :wait_for_ready, 2)
assert function_exported?(CodexSession, :send_turn, 2)
assert function_exported?(CodexSession, :interrupt_turn, 3)
assert function_exported?(CodexSession, :respond_to_approval, 3)
assert function_exported?(CodexSession, :respond_to_user_input, 3)
assert function_exported?(CodexSession, :read_thread, 2)
assert function_exported?(CodexSession, :rollback_thread, 3)
end
end
22 changes: 22 additions & 0 deletions apps/harness/test/harness/providers/cursor_session_test.exs
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
defmodule Harness.Providers.CursorSessionTest do
use ExUnit.Case, async: true

alias Harness.Providers.CursorSession

test "exports the provider session callbacks used by SessionManager" do
assert function_exported?(CursorSession, :start_link, 1)
assert function_exported?(CursorSession, :wait_for_ready, 1)
assert function_exported?(CursorSession, :wait_for_ready, 2)
assert function_exported?(CursorSession, :send_turn, 2)
assert function_exported?(CursorSession, :interrupt_turn, 3)
assert function_exported?(CursorSession, :respond_to_approval, 3)
assert function_exported?(CursorSession, :respond_to_user_input, 3)
assert function_exported?(CursorSession, :read_thread, 2)
assert function_exported?(CursorSession, :rollback_thread, 3)
end

test "reports rollback as unsupported" do
assert {:error, "Rollback not supported for Cursor provider"} =
CursorSession.rollback_thread(self(), "thread-1", 1)
end
end
17 changes: 17 additions & 0 deletions apps/harness/test/harness/providers/opencode_session_test.exs
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
defmodule Harness.Providers.OpenCodeSessionTest do
use ExUnit.Case, async: true

alias Harness.Providers.OpenCodeSession

test "exports the provider session callbacks used by SessionManager" do
assert function_exported?(OpenCodeSession, :start_link, 1)
assert function_exported?(OpenCodeSession, :wait_for_ready, 1)
assert function_exported?(OpenCodeSession, :wait_for_ready, 2)
assert function_exported?(OpenCodeSession, :send_turn, 2)
assert function_exported?(OpenCodeSession, :interrupt_turn, 3)
assert function_exported?(OpenCodeSession, :respond_to_approval, 3)
assert function_exported?(OpenCodeSession, :respond_to_user_input, 3)
assert function_exported?(OpenCodeSession, :read_thread, 2)
assert function_exported?(OpenCodeSession, :rollback_thread, 3)
end
end
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ import { ProjectionCheckpointRepository } from "../src/persistence/Services/Proj
import { ProjectionPendingApprovalRepository } from "../src/persistence/Services/ProjectionPendingApprovals.ts";
import { ProviderUnsupportedError } from "../src/provider/Errors.ts";
import { ProviderAdapterRegistry } from "../src/provider/Services/ProviderAdapterRegistry.ts";
import { McpConfigService } from "../src/provider/Services/McpConfig.ts";
import { ProviderSessionDirectoryLive } from "../src/provider/Layers/ProviderSessionDirectory.ts";
import { ServerSettingsService } from "../src/serverSettings.ts";
import { makeProviderServiceLive } from "../src/provider/Layers/ProviderService.ts";
Expand Down Expand Up @@ -270,6 +271,7 @@ export const makeOrchestrationIntegrationHarness = (
}),
).pipe(
Layer.provide(makeCodexAdapterLive()),
Layer.provideMerge(McpConfigService.layerTest()),
Layer.provideMerge(ServerConfig.layerTest(workspaceDir, rootDir)),
Layer.provideMerge(NodeServices.layer),
Layer.provideMerge(providerSessionDirectoryLayer),
Expand All @@ -279,11 +281,13 @@ export const makeOrchestrationIntegrationHarness = (
Layer.provide(providerSessionDirectoryLayer),
Layer.provide(realCodexRegistry),
Layer.provide(AnalyticsService.layerTest),
Layer.provide(McpConfigService.layerTest()),
)
: makeProviderServiceLive().pipe(
Layer.provide(providerSessionDirectoryLayer),
Layer.provide(fakeRegistry!),
Layer.provide(AnalyticsService.layerTest),
Layer.provide(McpConfigService.layerTest()),
);

const checkpointStoreLayer = CheckpointStoreLive.pipe(Layer.provide(GitCoreLive));
Expand Down
70 changes: 52 additions & 18 deletions apps/server/integration/TestProviderAdapter.integration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -240,6 +240,17 @@ export const makeTestProviderAdapterHarness = (options?: MakeTestProviderAdapter
>();

const emit = (event: ProviderRuntimeEvent) => Queue.offer(runtimeEvents, event);
const capabilities = {
sessionModelSwitch: provider === "cursor" ? "restart-session" : "in-session",
supportsUserInput: provider !== "cursor",
supportsRollback: provider !== "cursor",
supportsFileChangeApproval: provider !== "cursor",
resume: provider === "codex" ? "full" : "basic",
subagents: provider === "opencode" ? "full" : "none",
attachments: provider === "claudeAgent" ? "full" : "basic",
replay: "full",
mcpConfig: provider === "cursor" || provider === "claudeAgent" ? "none" : "basic",
} as const;

const startSession: ProviderAdapterShape<ProviderAdapterError>["startSession"] = (input) =>
Effect.gen(function* () {
Expand Down Expand Up @@ -401,23 +412,42 @@ export const makeTestProviderAdapterHarness = (options?: MakeTestProviderAdapter
requestId,
decision,
) =>
sessions.has(threadId)
? Effect.sync(() => {
const existing = approvalResponsesBySession.get(threadId) ?? [];
existing.push({
threadId,
requestId,
decision,
});
approvalResponsesBySession.set(threadId, existing);
})
: missingSessionEffect(provider, threadId);
!capabilities.supportsFileChangeApproval
? Effect.fail(
new ProviderAdapterValidationError({
provider,
operation: "respondToRequest",
issue: `${provider} does not support approval responses.`,
}),
)
: sessions.has(threadId)
? Effect.sync(() => {
const existing = approvalResponsesBySession.get(threadId) ?? [];
existing.push({
threadId,
requestId,
decision,
});
approvalResponsesBySession.set(threadId, existing);
})
: missingSessionEffect(provider, threadId);

const respondToUserInput: ProviderAdapterShape<ProviderAdapterError>["respondToUserInput"] = (
threadId,
_requestId,
_answers,
) => (sessions.has(threadId) ? Effect.void : missingSessionEffect(provider, threadId));
) =>
!capabilities.supportsUserInput
? Effect.fail(
new ProviderAdapterValidationError({
provider,
operation: "respondToUserInput",
issue: `${provider} does not support structured user input.`,
}),
)
: sessions.has(threadId)
? Effect.void
: missingSessionEffect(provider, threadId);

const stopSession: ProviderAdapterShape<ProviderAdapterError>["stopSession"] = (threadId) =>
Effect.sync(() => {
Expand All @@ -442,6 +472,15 @@ export const makeTestProviderAdapterHarness = (options?: MakeTestProviderAdapter
threadId,
numTurns,
) => {
if (!capabilities.supportsRollback) {
return Effect.fail(
new ProviderAdapterValidationError({
provider,
operation: "rollbackThread",
issue: `${provider} does not support rollback.`,
}),
);
}
const state = sessions.get(threadId);
if (!state) {
return missingSessionEffect(provider, threadId);
Expand Down Expand Up @@ -474,12 +513,7 @@ export const makeTestProviderAdapterHarness = (options?: MakeTestProviderAdapter

const adapter: ProviderAdapterShape<ProviderAdapterError> = {
provider,
capabilities: {
sessionModelSwitch: provider === "cursor" ? "restart-session" : "in-session",
supportsUserInput: provider !== "cursor",
supportsRollback: provider !== "cursor",
supportsFileChangeApproval: provider !== "cursor",
},
capabilities,
startSession,
sendTurn,
interruptTurn,
Expand Down
Loading
Loading