Skip to content
Draft
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
53 changes: 53 additions & 0 deletions clients/dotnet/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -50,9 +50,62 @@ var state = new SessionState { /* ... */ };
Reducers.ApplyToSession(state, action); // mutates `state` in place
```

For TCP accounting, `Reducers.TcpReducer(state, action)` returns a new state
and throws `InvalidOperationException` on invalid actions.

See [`examples/`](examples/) for runnable `ConnectWs` and `ReducersDemo`
console apps.

## TCP streams

After `InitializeAsync` negotiates `TcpConnections`, the concrete `AhpClient`
provides an integrated adapter:

```csharp
await using var tcp = await client.OpenTcpConnectionAsync(sessionUri,
new TcpConnectionSubscription
{
Type = "tcpConnection", Host = "localhost", Port = 3000,
Encoding = TcpDataEncoding.Base64,
ReceiveWindowBytes = 65536, MaximumChunkSize = 16384,
}, cancellationToken);
await tcp.WriteAsync(requestBytes, cancellationToken);
await tcp.EndAsync(cancellationToken); // input EOF; output remains readable
while (await tcp.ReadAsync(cancellationToken) is { } bytes)
await ConsumeAsync(bytes);
```

The SDK owns buffering, flow control, and replay. One reader and one writer may
run concurrently. `ReadAsync` releases receive credit; `DrainAsync` waits for
destination consumption. `CloseAsync` stops writes but retains crossing output,
so keep reading during graceful close. Disposal aborts without draining.

Transport loss suspends the same handles. On a fresh `AhpClient`, call
`ReconnectTcpConnectionsAsync(reconnectParams, handles, token)` with the original
`ClientId`; apply its returned replay to ordinary subscriptions. For deliberate
transport replacement, use `ShutdownAsync(preserveTcpConnections: true)` instead
of normal shutdown, which terminates streams.

With `MultiHostClient`, use `HostClientHandle.OpenTcpConnectionAsync(session, create)`
for automatic reconnect. Host removal/shutdown terminates retained streams.
Missing resources and snapshot fallback fail streams rather than creating new
sockets. TCP streams do not belong in ordinary state mirrors.

Transport policy, connection limits, and native socket bridges remain
application-owned. See the [TCP channel contract](../../docs/specification/tcp-channel.md).

## Strict events

For custom loss-sensitive consumers, attach
`client.CreateEventStream(failOnOverflow: true)` before sending requests and
retain it for `Events.ReadAllAsync`. Overflow throws `SubscriptionLagException`;
decode loss throws `AhpTransportException` of kind `"protocol"`. Both terminate
the receiver rather than skipping events. Capacity uses
`ClientConfig.SubscriptionBufferCapacity`; ordinary receivers are unchanged.
These raw receivers are global. The owned TCP adapter instead registers a strict
child-scoped receiver during creation reply processing and reattaches it per
child on reconnect. Unrelated traffic cannot exhaust a TCP stream's event buffer.

## Dependency injection

Register the services with `AddAgentHostProtocol` (in the
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -231,6 +231,26 @@ public ActionType(string value)

public static readonly ActionType AutomationRunCancelRequested = new ActionType("automationRun/cancelRequested");

public static readonly ActionType TcpInput = new ActionType("tcp/input");

public static readonly ActionType TcpData = new ActionType("tcp/data");

public static readonly ActionType TcpInputConsumed = new ActionType("tcp/inputConsumed");

public static readonly ActionType TcpDataConsumed = new ActionType("tcp/dataConsumed");

public static readonly ActionType TcpInputEof = new ActionType("tcp/inputEof");

public static readonly ActionType TcpDataEof = new ActionType("tcp/dataEof");

public static readonly ActionType TcpClientClose = new ActionType("tcp/clientClose");

public static readonly ActionType TcpHostClose = new ActionType("tcp/hostClose");

public static readonly ActionType TcpClientReset = new ActionType("tcp/clientReset");

public static readonly ActionType TcpHostReset = new ActionType("tcp/hostReset");

/// <inheritdoc />
public bool Equals(ActionType other) => string.Equals(Value, other.Value, StringComparison.Ordinal);

Expand Down Expand Up @@ -2400,6 +2420,92 @@ public sealed record ResourceWatchChangedAction
public JsonElement Changes { get; init; }
}

/// <summary>Client bytes. Never apply optimistically to the authoritative reducer.
/// Write to the destination only when accepted input.receivedBytes advances.</summary>
public sealed record TcpInputAction
{
public ActionType Type { get; init; } = ActionType.TcpInput;

/// <summary>Absolute decoded-byte offset.</summary>
public long Offset { get; init; }

/// <summary>Nonempty canonical padded RFC 4648 base64; no whitespace.</summary>
public required string Data { get; init; }
}

/// <summary>Host bytes. Deliver once, only when output.receivedBytes advances.</summary>
public sealed record TcpDataAction
{
public ActionType Type { get; init; } = ActionType.TcpData;

/// <summary>Absolute decoded-byte offset.</summary>
public long Offset { get; init; }

/// <summary>Nonempty canonical padded RFC 4648 base64; no whitespace.</summary>
public required string Data { get; init; }
}

/// <summary>Cumulative input bytes released from the host's bounded write buffer.
/// Not an acknowledgment that the destination application processed the bytes.</summary>
public sealed record TcpInputConsumedAction
{
public ActionType Type { get; init; } = ActionType.TcpInputConsumed;

public long ConsumedBytes { get; init; }
}

/// <summary>Cumulative output bytes released by the client's bounded stream consumer.</summary>
public sealed record TcpDataConsumedAction
{
public ActionType Type { get; init; } = ActionType.TcpDataConsumed;

public long ConsumedBytes { get; init; }
}

/// <summary>Half-close client input after all preceding input bytes have been written.</summary>
public sealed record TcpInputEofAction
{
public ActionType Type { get; init; } = ActionType.TcpInputEof;

public long FinalOffset { get; init; }
}

/// <summary>Half-close host output after all preceding output bytes have been delivered.</summary>
public sealed record TcpDataEofAction
{
public ActionType Type { get; init; } = ActionType.TcpDataEof;

public long FinalOffset { get; init; }
}

/// <summary>Client's final close. Respond with hostClose if not already sent.</summary>
public sealed record TcpClientCloseAction
{
public ActionType Type { get; init; } = ActionType.TcpClientClose;
}

/// <summary>Host's final close. Respond with clientClose if not already sent.</summary>
public sealed record TcpHostCloseAction
{
public ActionType Type { get; init; } = ActionType.TcpHostClose;
}

/// <summary>Abort both directions and discard buffered payload.</summary>
public sealed record TcpClientResetAction
{
public ActionType Type { get; init; } = ActionType.TcpClientReset;

public TcpResetReason Reason { get; init; }
}

/// <summary>Abort both directions and discard buffered payload.</summary>
public sealed record TcpHostResetAction
{
public ActionType Type { get; init; } = ActionType.TcpHostReset;

public TcpResetReason Reason { get; init; }
}

/// <summary>Upsert an {@link Annotation} in the annotations channel — adds a new
/// annotation, or replaces an existing one identified by
/// {@link Annotation.id}.
Expand Down Expand Up @@ -2836,6 +2942,16 @@ public StateActionConverter()
["terminal/commandExecuted"] = typeof(TerminalCommandExecutedAction),
["terminal/commandFinished"] = typeof(TerminalCommandFinishedAction),
["resourceWatch/changed"] = typeof(ResourceWatchChangedAction),
["tcp/input"] = typeof(TcpInputAction),
["tcp/data"] = typeof(TcpDataAction),
["tcp/inputConsumed"] = typeof(TcpInputConsumedAction),
["tcp/dataConsumed"] = typeof(TcpDataConsumedAction),
["tcp/inputEof"] = typeof(TcpInputEofAction),
["tcp/dataEof"] = typeof(TcpDataEofAction),
["tcp/clientClose"] = typeof(TcpClientCloseAction),
["tcp/hostClose"] = typeof(TcpHostCloseAction),
["tcp/clientReset"] = typeof(TcpClientResetAction),
["tcp/hostReset"] = typeof(TcpHostResetAction),
["annotations/set"] = typeof(AnnotationsSetAction),
["annotations/removed"] = typeof(AnnotationsRemovedAction),
["annotations/entrySet"] = typeof(AnnotationsEntrySetAction),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -394,6 +394,10 @@ public sealed record InitializeResult
/// host does not expose an automation catalogue or automation commands.</summary>
[JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)]
public AutomationCapabilities? Automations { get; init; }

/// <summary>Enables atomic creation of session-scoped, replay-only TCP channels.</summary>
[JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)]
public TcpConnectionsCapability? TcpConnections { get; init; }
}

/// <summary>Identifies a protocol implementation — the software (and build) on one end
Expand Down Expand Up @@ -571,6 +575,12 @@ public sealed record ReconnectSnapshotResult

/// <summary>Fresh snapshots for each subscription</summary>
public required List<Snapshot> Snapshots { get; init; }

/// <summary>Subscriptions that cannot be restored. Hosts supporting TCP MUST list all
/// requested TCP channels here and dispose their sockets on snapshot fallback.
/// Omitted by older hosts; absence does not authorize snapshot-restoring TCP.</summary>
[JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)]
public List<string>? Missing { get; init; }
}

/// <summary>Subscribe to a URI-identified channel.
Expand Down Expand Up @@ -604,6 +614,12 @@ public sealed record SubscribeParams
/// default snapshot. Clients MUST tolerate receiving more state than requested.</summary>
[JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)]
public SubscribeView? View { get; init; }

/// <summary>Atomically create a private child channel and subscribe to it.
/// Requires the advertised tcpConnections capability. channel identifies
/// the parent session; snapshot.resource identifies the created TCP channel.</summary>
[JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)]
public TcpConnectionSubscription? Create { get; init; }
}

/// <summary>Optional client-requested shape for a subscription snapshot.</summary>
Expand Down Expand Up @@ -643,6 +659,32 @@ public sealed record SubscribeResult
public Snapshot? Snapshot { get; init; }
}

/// <summary>Creates and exclusively subscribes to one TCP connection.
///
/// SubscribeParams.channel MUST identify the parent `ahp-session:` channel.
/// The host returns the new `ahp-tcp:` URI in snapshot.resource, not the parent.
/// It installs the subscription and sends the response before any TCP actions.
/// Unknown creation kinds MUST be rejected, never treated as normal subscribe.</summary>
public sealed record TcpConnectionSubscription
{
public required string Type { get; init; }

/// <summary>DNS name or IP literal, not a URL.</summary>
public required string Host { get; init; }

/// <summary>Destination port.</summary>
public long Port { get; init; }

/// <summary>Selected from InitializeResult.tcpConnections.encodings.</summary>
public TcpDataEncoding Encoding { get; init; }

/// <summary>Client receive window in decoded bytes.</summary>
public long ReceiveWindowBytes { get; init; }

/// <summary>Maximum decoded bytes per output action; MUST NOT exceed receiveWindowBytes.</summary>
public long MaximumChunkSize { get; init; }
}

// TODO: could not generate SessionForkSource: Error: Interface SessionForkSource not found

/// <summary>Creates a new session with the specified agent provider.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,8 @@ public static class AhpErrorCodes
public const int AlreadyExists = -32010;
/// <summary>An optimistic-concurrency precondition failed. Returned when a request carries a precondition token that no longer matches the receiver's current state — for example, `resourceWrite` with an `ifMatch` etag that has been superseded by a concurrent write. Callers SHOULD re-read the resource (e.g. via `resourceResolve`) and decide whether to retry the operation with the fresh token or surface the conflict to the user.</summary>
public const int Conflict = -32011;
/// <summary>TCP creation failed; data MUST contain TcpConnectionOpenErrorData.</summary>
public const int TcpConnectionOpenFailed = -32012;
}

/// <summary>Detail payload of an AuthRequired (-32007) error.</summary>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -221,6 +221,7 @@ namespace Microsoft.AgentHostProtocol;
[JsonSerializable(typeof(FileEditCollection))]
[JsonSerializable(typeof(FileEditDiffStats))]
[JsonSerializable(typeof(FileEditSide))]
[JsonSerializable(typeof(FlowControlledByteDirectionState))]
[JsonSerializable(typeof(ForkChatSource))]
[JsonSerializable(typeof(HookCustomization))]
[JsonSerializable(typeof(Icon))]
Expand Down Expand Up @@ -411,6 +412,26 @@ namespace Microsoft.AgentHostProtocol;
[JsonSerializable(typeof(SubscribeView))]
[JsonSerializable(typeof(SubscriptionDeliveryOptions))]
[JsonSerializable(typeof(SystemNotificationResponsePart))]
[JsonSerializable(typeof(TcpClientCloseAction))]
[JsonSerializable(typeof(TcpClientResetAction))]
[JsonSerializable(typeof(TcpConnectionOpenErrorData))]
[JsonSerializable(typeof(TcpConnectionOpenFailureReason))]
[JsonSerializable(typeof(TcpConnectionsCapability))]
[JsonSerializable(typeof(TcpConnectionState))]
[JsonSerializable(typeof(TcpConnectionSubscription))]
[JsonSerializable(typeof(TcpDataAction))]
[JsonSerializable(typeof(TcpDataConsumedAction))]
[JsonSerializable(typeof(TcpDataEncoding))]
[JsonSerializable(typeof(TcpDataEofAction))]
[JsonSerializable(typeof(TcpEndpoint))]
[JsonSerializable(typeof(TcpHostCloseAction))]
[JsonSerializable(typeof(TcpHostResetAction))]
[JsonSerializable(typeof(TcpInputAction))]
[JsonSerializable(typeof(TcpInputConsumedAction))]
[JsonSerializable(typeof(TcpInputEofAction))]
[JsonSerializable(typeof(TcpResetReason))]
[JsonSerializable(typeof(TcpResetState))]
[JsonSerializable(typeof(TcpTarget))]
[JsonSerializable(typeof(TelemetryCapabilities))]
[JsonSerializable(typeof(TerminalClaim))]
[JsonSerializable(typeof(TerminalClaimedAction))]
Expand Down
Loading
Loading