/*---------------------------------------------------------------------------------------------
* Copyright (c) Microsoft Corporation. All rights reserved.
*--------------------------------------------------------------------------------------------*/
using GitHub.Copilot.SDK.Rpc;
using Microsoft.Extensions.AI;
using Microsoft.Extensions.Logging;
using StreamJsonRpc;
using System.Collections.Immutable;
using System.Text.Json;
using System.Text.Json.Nodes;
using System.Text.Json.Serialization;
using System.Threading.Channels;
namespace GitHub.Copilot.SDK;
///
/// Represents a single conversation session with the Copilot CLI.
///
///
///
/// A session maintains conversation state, handles events, and manages tool execution.
/// Sessions are created via or resumed via
/// .
///
///
/// The session provides methods to send messages, subscribe to events, retrieve
/// conversation history, and manage the session lifecycle.
///
///
/// implements . Use the
/// await using pattern for automatic cleanup, or call
/// explicitly. Disposing a session releases in-memory resources but preserves session data
/// on disk — the conversation can be resumed later via
/// . To permanently delete session data,
/// use .
///
///
///
///
/// await using var session = await client.CreateSessionAsync(new() { OnPermissionRequest = PermissionHandler.ApproveAll, Model = "gpt-4" });
///
/// // Subscribe to events
/// using var subscription = session.On(evt =>
/// {
/// if (evt is AssistantMessageEvent assistantMessage)
/// {
/// Console.WriteLine($"Assistant: {assistantMessage.Data?.Content}");
/// }
/// });
///
/// // Send a message and wait for completion
/// await session.SendAndWaitAsync(new MessageOptions { Prompt = "Hello, world!" });
///
///
public sealed partial class CopilotSession : IAsyncDisposable
{
private readonly Dictionary _toolHandlers = [];
private readonly JsonRpc _rpc;
private readonly ILogger _logger;
private volatile PermissionRequestHandler? _permissionHandler;
private volatile UserInputHandler? _userInputHandler;
private ImmutableArray _eventHandlers = ImmutableArray.Empty;
private SessionHooks? _hooks;
private readonly SemaphoreSlim _hooksLock = new(1, 1);
private SessionRpc? _sessionRpc;
private int _isDisposed;
///
/// Channel that serializes event dispatch. enqueues;
/// a single background consumer () dequeues and
/// invokes handlers one at a time, preserving arrival order.
///
private readonly Channel _eventChannel = Channel.CreateUnbounded(
new() { SingleReader = true });
///
/// Gets the unique identifier for this session.
///
/// A string that uniquely identifies this session.
public string SessionId { get; }
///
/// Gets the typed RPC client for session-scoped methods.
///
public SessionRpc Rpc => _sessionRpc ??= new SessionRpc(_rpc, SessionId);
///
/// Gets the path to the session workspace directory when infinite sessions are enabled.
///
///
/// The path to the workspace containing checkpoints/, plan.md, and files/ subdirectories,
/// or null if infinite sessions are disabled.
///
public string? WorkspacePath { get; internal set; }
///
/// Initializes a new instance of the class.
///
/// The unique identifier for this session.
/// The JSON-RPC connection to the Copilot CLI.
/// Logger for diagnostics.
/// The workspace path if infinite sessions are enabled.
///
/// This constructor is internal. Use to create sessions.
///
internal CopilotSession(string sessionId, JsonRpc rpc, ILogger logger, string? workspacePath = null)
{
SessionId = sessionId;
_rpc = rpc;
_logger = logger;
WorkspacePath = workspacePath;
// Start the asynchronous processing loop.
_ = ProcessEventsAsync();
}
private Task InvokeRpcAsync(string method, object?[]? args, CancellationToken cancellationToken)
{
return CopilotClient.InvokeRpcAsync(_rpc, method, args, cancellationToken);
}
///
/// Sends a message to the Copilot session and waits for the response.
///
/// Options for the message to be sent, including the prompt and optional attachments.
/// A that can be used to cancel the operation.
/// A task that resolves with the ID of the response message, which can be used to correlate events.
/// Thrown if the session has been disposed.
///
///
/// This method returns immediately after the message is queued. Use
/// if you need to wait for the assistant to finish processing.
///
///
/// Subscribe to events via to receive streaming responses and other session events.
///
///
///
///
/// var messageId = await session.SendAsync(new MessageOptions
/// {
/// Prompt = "Explain this code",
/// Attachments = new List<Attachment>
/// {
/// new() { Type = "file", Path = "./Program.cs" }
/// }
/// });
///
///
public async Task SendAsync(MessageOptions options, CancellationToken cancellationToken = default)
{
var (traceparent, tracestate) = TelemetryHelpers.GetTraceContext();
var request = new SendMessageRequest
{
SessionId = SessionId,
Prompt = options.Prompt,
Attachments = options.Attachments,
Mode = options.Mode,
Traceparent = traceparent,
Tracestate = tracestate
};
var response = await InvokeRpcAsync(
"session.send", [request], cancellationToken);
return response.MessageId;
}
///
/// Sends a message to the Copilot session and waits until the session becomes idle.
///
/// Options for the message to be sent, including the prompt and optional attachments.
/// Timeout duration (default: 60 seconds). Controls how long to wait; does not abort in-flight agent work.
/// A that can be used to cancel the operation.
/// A task that resolves with the final assistant message event, or null if none was received.
/// Thrown if the timeout is reached before the session becomes idle.
/// Thrown if the is cancelled.
/// Thrown if the session has been disposed.
///
///
/// This is a convenience method that combines with waiting for
/// the session.idle event. Use this when you want to block until the assistant
/// has finished processing the message.
///
///
/// Events are still delivered to handlers registered via while waiting.
///
///
///
///
/// // Send and wait for completion with default 60s timeout
/// var response = await session.SendAndWaitAsync(new MessageOptions { Prompt = "What is 2+2?" });
/// Console.WriteLine(response?.Data?.Content); // "4"
///
///
public async Task SendAndWaitAsync(
MessageOptions options,
TimeSpan? timeout = null,
CancellationToken cancellationToken = default)
{
var effectiveTimeout = timeout ?? TimeSpan.FromSeconds(60);
var tcs = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
AssistantMessageEvent? lastAssistantMessage = null;
void Handler(SessionEvent evt)
{
switch (evt)
{
case AssistantMessageEvent assistantMessage:
lastAssistantMessage = assistantMessage;
break;
case SessionIdleEvent:
tcs.TrySetResult(lastAssistantMessage);
break;
case SessionErrorEvent errorEvent:
var message = errorEvent.Data?.Message ?? "session error";
tcs.TrySetException(new InvalidOperationException($"Session error: {message}"));
break;
}
}
using var subscription = On(Handler);
await SendAsync(options, cancellationToken);
using var cts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
cts.CancelAfter(effectiveTimeout);
using var registration = cts.Token.Register(() =>
{
if (cancellationToken.IsCancellationRequested)
tcs.TrySetCanceled(cancellationToken);
else
tcs.TrySetException(new TimeoutException($"SendAndWaitAsync timed out after {effectiveTimeout}"));
});
return await tcs.Task;
}
///
/// Registers a callback for session events.
///
/// A callback to be invoked when a session event occurs.
/// An that, when disposed, unsubscribes the handler.
///
///
/// Events include assistant messages, tool executions, errors, and session state changes.
/// Multiple handlers can be registered and will all receive events.
///
///
/// Handlers are invoked serially in event-arrival order on a background thread.
/// A handler will never be called concurrently with itself or with other handlers
/// on the same session.
///
///
///
///
/// using var subscription = session.On(evt =>
/// {
/// switch (evt)
/// {
/// case AssistantMessageEvent:
/// Console.WriteLine($"Assistant: {evt.Data?.Content}");
/// break;
/// case SessionErrorEvent:
/// Console.WriteLine($"Error: {evt.Data?.Message}");
/// break;
/// }
/// });
///
/// // The handler is automatically unsubscribed when the subscription is disposed.
///
///
public IDisposable On(SessionEventHandler handler)
{
ImmutableInterlocked.Update(ref _eventHandlers, array => array.Add(handler));
return new ActionDisposable(() => ImmutableInterlocked.Update(ref _eventHandlers, array => array.Remove(handler)));
}
///
/// Enqueues an event for serial dispatch to all registered handlers.
///
/// The session event to dispatch.
///
/// This method is non-blocking. Broadcast request events (external_tool.requested,
/// permission.requested) are fired concurrently so that a stalled handler does not
/// block event delivery. The event is then placed into an in-memory channel and
/// processed by a single background consumer (),
/// which guarantees user handlers see events one at a time, in order.
///
internal void DispatchEvent(SessionEvent sessionEvent)
{
// Fire broadcast work concurrently (fire-and-forget with error logging).
// This is done outside the channel so broadcast handlers don't block the
// consumer loop — important when a secondary client's handler intentionally
// never completes (multi-client permission scenario).
_ = HandleBroadcastEventAsync(sessionEvent);
// Queue the event for serial processing by user handlers.
_eventChannel.Writer.TryWrite(sessionEvent);
}
///
/// Single-reader consumer loop that processes events from the channel.
/// Ensures user event handlers are invoked serially and in FIFO order.
///
private async Task ProcessEventsAsync()
{
await foreach (var sessionEvent in _eventChannel.Reader.ReadAllAsync())
{
foreach (var handler in _eventHandlers)
{
try
{
handler(sessionEvent);
}
catch (Exception ex)
{
LogEventHandlerError(ex);
}
}
}
}
///
/// Registers custom tool handlers for this session.
///
/// A collection of AI functions that can be invoked by the assistant.
///
/// Tools allow the assistant to execute custom functions. When the assistant invokes a tool,
/// the corresponding handler is called with the tool arguments.
///
internal void RegisterTools(ICollection tools)
{
_toolHandlers.Clear();
foreach (var tool in tools)
{
_toolHandlers.Add(tool.Name, tool);
}
}
///
/// Retrieves a registered tool by name.
///
/// The name of the tool to retrieve.
/// The tool if found; otherwise, null.
internal AIFunction? GetTool(string name)
{
return _toolHandlers.TryGetValue(name, out var tool) ? tool : null;
}
///
/// Registers a handler for permission requests.
///
/// The permission handler function.
///
/// When the assistant needs permission to perform certain actions (e.g., file operations),
/// this handler is called to approve or deny the request.
///
internal void RegisterPermissionHandler(PermissionRequestHandler handler)
{
_permissionHandler = handler;
}
///
/// Handles a permission request from the Copilot CLI.
///
/// The permission request data from the CLI.
/// A task that resolves with the permission decision.
internal async Task HandlePermissionRequestAsync(JsonElement permissionRequestData)
{
var handler = _permissionHandler;
if (handler == null)
{
return new PermissionRequestResult
{
Kind = PermissionRequestResultKind.DeniedCouldNotRequestFromUser
};
}
var request = JsonSerializer.Deserialize(permissionRequestData.GetRawText(), SessionEventsJsonContext.Default.PermissionRequest)
?? throw new InvalidOperationException("Failed to deserialize permission request");
var invocation = new PermissionInvocation
{
SessionId = SessionId
};
return await handler(request, invocation);
}
///
/// Handles broadcast request events by executing local handlers and responding via RPC.
/// Implements the protocol v3 broadcast model where tool calls and permission requests
/// are broadcast as session events to all clients.
///
private async Task HandleBroadcastEventAsync(SessionEvent sessionEvent)
{
try
{
switch (sessionEvent)
{
case ExternalToolRequestedEvent toolEvent:
{
var data = toolEvent.Data;
if (string.IsNullOrEmpty(data.RequestId) || string.IsNullOrEmpty(data.ToolName))
return;
var tool = GetTool(data.ToolName);
if (tool is null)
return; // This client doesn't handle this tool; another client will.
using (TelemetryHelpers.RestoreTraceContext(data.Traceparent, data.Tracestate))
await ExecuteToolAndRespondAsync(data.RequestId, data.ToolName, data.ToolCallId, data.Arguments, tool);
break;
}
case PermissionRequestedEvent permEvent:
{
var data = permEvent.Data;
if (string.IsNullOrEmpty(data.RequestId) || data.PermissionRequest is null)
return;
var handler = _permissionHandler;
if (handler is null)
return; // This client doesn't handle permissions; another client will.
await ExecutePermissionAndRespondAsync(data.RequestId, data.PermissionRequest, handler);
break;
}
}
}
catch (Exception ex) when (ex is not OperationCanceledException)
{
LogBroadcastHandlerError(ex);
}
}
///
/// Executes a tool handler and sends the result back via the HandlePendingToolCall RPC.
///
private async Task ExecuteToolAndRespondAsync(string requestId, string toolName, string toolCallId, object? arguments, AIFunction tool)
{
try
{
var invocation = new ToolInvocation
{
SessionId = SessionId,
ToolCallId = toolCallId,
ToolName = toolName,
Arguments = arguments
};
var aiFunctionArgs = new AIFunctionArguments
{
Context = new Dictionary