Skip to content

Commit b62055d

Browse files
[http-client-csharp] Add JSONL and SSE streaming support (#11462)
## Summary - preserve JSONL and SSE stream metadata in the C# emitter code model - generate incremental `IAsyncEnumerable<T>` JSONL requests and `AsyncStreamingClientResult<T>` JSONL responses - generate unbuffered SSE responses as `AsyncStreamingClientResult<SseItem<BinaryData>>`, including terminal-event suppression - add live Spector coverage for JSONL send/receive and all unnamed, named, and retrieval SSE scenarios - consume the streaming result APIs introduced by Azure/azure-sdk-for-net#60994 Closes #10738 ## Validation - `npm run build` - `npm test` - `npm run cop` - `Test-Spector.ps1 -filter "http/streaming/jsonl"` - `Test-Spector.ps1 -filter "http/streaming/sse"` `npm run lint -- --emitter` is currently blocked by the package resolving ESLint 10 while retaining a legacy ESLint configuration (`eslint.config.*` is not present). --------- Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: e02ee7d0-8d27-46f8-a856-222d6a938ecd
1 parent 0fba3fe commit b62055d

141 files changed

Lines changed: 3948 additions & 1049 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

packages/http-client-csharp/emitter/src/lib/operation-converter.ts

Lines changed: 112 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ import type {
1919
SdkQueryParameter,
2020
SdkServiceMethod,
2121
SdkServiceResponseHeader,
22+
SdkStreamMetadata,
2223
SdkType,
2324
} from "@azure-tools/typespec-client-generator-core";
2425
import {
@@ -30,15 +31,18 @@ import {
3031
shouldGenerateConvenient,
3132
shouldGenerateProtocol,
3233
} from "@azure-tools/typespec-client-generator-core";
33-
import type { Diagnostic } from "@typespec/compiler";
34+
import type { Diagnostic, Union } from "@typespec/compiler";
3435
import {
36+
compilerAssert,
3537
createDiagnosticCollector,
3638
getDeprecated,
3739
isErrorModel,
3840
NoTarget,
3941
} from "@typespec/compiler";
42+
import { unsafe_getEventDefinitions } from "@typespec/events/experimental";
4043
import type { HttpStatusCodeRange } from "@typespec/http";
4144
import { getResourceOperation } from "@typespec/rest";
45+
import { isTerminalEvent } from "@typespec/sse";
4246
import type { CSharpEmitterContext } from "../sdk-context.js";
4347
import { collectionFormatToDelimMap } from "../type/collection-format.js";
4448
import type { HttpResponseHeader } from "../type/http-response-header.js";
@@ -64,6 +68,7 @@ import type {
6468
InputMethodParameter,
6569
InputPathParameter,
6670
InputQueryParameter,
71+
InputStreamingType,
6772
InputType,
6873
} from "../type/input-type.js";
6974
import { convertLroFinalStateVia } from "../type/operation-final-state-via.js";
@@ -234,7 +239,9 @@ export function fromSdkServiceMethodOperation(
234239
path: method.operation.path,
235240
externalDocsUrl: getExternalDocs(sdkContext, method.operation.__raw.operation)?.url,
236241
requestMediaTypes: requestMediaTypes,
237-
bufferResponse: true,
242+
bufferResponse: !method.operation.responses.some((response) =>
243+
isSupportedStream(response.streamMetadata),
244+
),
238245
generateProtocolMethod: shouldGenerateProtocol(sdkContext, method.operation.__raw.operation),
239246
generateConvenienceMethod: generateConvenience,
240247
crossLanguageDefinitionId: method.crossLanguageDefinitionId,
@@ -351,8 +358,23 @@ function fromSdkServiceMethodParameters(
351358
for (const p of method.parameters) {
352359
const methodInputParameter = diagnostics.pipe(fromMethodParameter(sdkContext, p, namespace));
353360
const operationHttpParameter = getHttpOperationParameter(method, p);
361+
const streamMetadata = method.operation.bodyParam?.streamMetadata;
362+
const isStreamingBodyParameter =
363+
isJsonLinesStream(streamMetadata) &&
364+
method.operation.bodyParam!.methodParameterSegments.some((segments) =>
365+
segments.some(
366+
(segment) =>
367+
segment === p || segment.crossLanguageDefinitionId === p.crossLanguageDefinitionId,
368+
),
369+
);
354370

355371
if (!operationHttpParameter) {
372+
if (isStreamingBodyParameter) {
373+
methodInputParameter.type = diagnostics.pipe(
374+
fromSdkStreamMetadata(sdkContext, streamMetadata),
375+
);
376+
methodInputParameter.location = RequestLocation.Body;
377+
}
356378
parameters.push(methodInputParameter);
357379
continue;
358380
}
@@ -365,6 +387,12 @@ function fromSdkServiceMethodParameters(
365387
rootApiVersions,
366388
diagnostics,
367389
);
390+
if (isStreamingBodyParameter) {
391+
methodInputParameter.type = diagnostics.pipe(
392+
fromSdkStreamMetadata(sdkContext, streamMetadata),
393+
);
394+
methodInputParameter.location = RequestLocation.Body;
395+
}
368396
parameters.push(methodInputParameter);
369397
}
370398

@@ -378,6 +406,15 @@ function updateMethodParameter(
378406
rootApiVersions: string[],
379407
diagnostics: ReturnType<typeof createDiagnosticCollector>,
380408
): void {
409+
if (
410+
operationHttpParameter.kind === "body" &&
411+
isJsonLinesStream(operationHttpParameter.streamMetadata)
412+
) {
413+
methodParameter.type = diagnostics.pipe(
414+
fromSdkStreamMetadata(sdkContext, operationHttpParameter.streamMetadata),
415+
);
416+
}
417+
381418
// for content type parameter
382419
if (isContentType(operationHttpParameter)) {
383420
methodParameter.type = diagnostics.pipe(
@@ -408,7 +445,9 @@ function fromSdkServiceMethodResponse(
408445
const diagnostics = createDiagnosticCollector();
409446

410447
return diagnostics.wrap({
411-
type: diagnostics.pipe(getResponseType(sdkContext, methodResponse.type)),
448+
type: isSupportedStream(methodResponse.streamMetadata)
449+
? diagnostics.pipe(fromSdkStreamMetadata(sdkContext, methodResponse.streamMetadata))
450+
: diagnostics.pipe(getResponseType(sdkContext, methodResponse.type)),
412451
resultSegments: methodResponse.resultSegments?.map((segment) =>
413452
getResponseSegmentName(segment),
414453
),
@@ -607,7 +646,9 @@ function fromBodyParameter(
607646
rootApiVersions: string[],
608647
): [InputBodyParameter, readonly Diagnostic[]] {
609648
const diagnostics = createDiagnosticCollector();
610-
const parameterType = diagnostics.pipe(fromSdkType(sdkContext, p.type, p));
649+
const parameterType = isJsonLinesStream(p.streamMetadata)
650+
? diagnostics.pipe(fromSdkStreamMetadata(sdkContext, p.streamMetadata))
651+
: diagnostics.pipe(fromSdkType(sdkContext, p.type, p));
611652

612653
const retVar: InputBodyParameter = {
613654
kind: "body",
@@ -725,7 +766,9 @@ export function fromSdkHttpOperationResponse(
725766
const range = sdkResponse.statusCodes;
726767
retVar = {
727768
statusCodes: toStatusCodesArray(range),
728-
bodyType: diagnostics.pipe(getResponseType(sdkContext, sdkResponse.type)),
769+
bodyType: isSupportedStream(sdkResponse.streamMetadata)
770+
? diagnostics.pipe(fromSdkStreamMetadata(sdkContext, sdkResponse.streamMetadata))
771+
: diagnostics.pipe(getResponseType(sdkContext, sdkResponse.type)),
729772
headers: diagnostics.pipe(fromSdkServiceResponseHeaders(sdkContext, sdkResponse.headers)),
730773
isErrorResponse:
731774
sdkResponse.type !== undefined && isErrorModel(sdkContext.program, sdkResponse.type.__raw!),
@@ -737,6 +780,70 @@ export function fromSdkHttpOperationResponse(
737780
return diagnostics.wrap(retVar);
738781
}
739782

783+
function fromSdkStreamMetadata(
784+
sdkContext: CSharpEmitterContext,
785+
streamMetadata: SdkStreamMetadata,
786+
): [InputStreamingType, readonly Diagnostic[]] {
787+
const diagnostics = createDiagnosticCollector();
788+
const originalType = streamMetadata.originalType;
789+
const streamKind = getStreamKind(streamMetadata);
790+
compilerAssert(
791+
streamKind !== undefined,
792+
"Stream metadata must have a supported stream kind before it is converted.",
793+
);
794+
let terminalEventType: string | undefined;
795+
let terminalEventValue: string | undefined;
796+
797+
if (streamKind === "sse" && streamMetadata.streamType.__raw?.kind === "Union") {
798+
const eventDefinitions = diagnostics.pipe(
799+
unsafe_getEventDefinitions(sdkContext.program, streamMetadata.streamType.__raw as Union),
800+
);
801+
const terminalDefinition = eventDefinitions.find((definition) =>
802+
isTerminalEvent(sdkContext.program, definition.root),
803+
);
804+
if (terminalDefinition?.payloadType.kind === "String") {
805+
terminalEventType = terminalDefinition.eventType;
806+
terminalEventValue = terminalDefinition.payloadType.value;
807+
}
808+
}
809+
810+
return diagnostics.wrap({
811+
kind: "streaming",
812+
name: originalType.kind === "model" ? originalType.name : "Stream",
813+
valueType: diagnostics.pipe(fromSdkType(sdkContext, streamMetadata.streamType)),
814+
streamKind,
815+
contentTypes: streamMetadata.contentTypes,
816+
terminalEventType,
817+
terminalEventValue,
818+
crossLanguageDefinitionId:
819+
originalType.kind === "model" ? originalType.crossLanguageDefinitionId : "",
820+
});
821+
}
822+
823+
function getStreamKind(
824+
streamMetadata: SdkStreamMetadata | undefined,
825+
): InputStreamingType["streamKind"] | undefined {
826+
if (streamMetadata?.contentTypes.includes("application/jsonl")) {
827+
return "jsonl";
828+
}
829+
if (streamMetadata?.contentTypes.includes("text/event-stream")) {
830+
return "sse";
831+
}
832+
return undefined;
833+
}
834+
835+
function isSupportedStream(
836+
streamMetadata: SdkStreamMetadata | undefined,
837+
): streamMetadata is SdkStreamMetadata {
838+
return getStreamKind(streamMetadata) !== undefined;
839+
}
840+
841+
function isJsonLinesStream(
842+
streamMetadata: SdkStreamMetadata | undefined,
843+
): streamMetadata is SdkStreamMetadata {
844+
return getStreamKind(streamMetadata) === "jsonl";
845+
}
846+
740847
function fromSdkServiceResponseHeaders(
741848
sdkContext: CSharpEmitterContext,
742849
headers: SdkServiceResponseHeader[],

packages/http-client-csharp/emitter/src/type/input-type.ts

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -79,6 +79,7 @@ export type InputType =
7979
| InputEnumType
8080
| InputEnumValueType
8181
| InputArrayType
82+
| InputStreamingType
8283
| InputDictionaryType
8384
| InputNullableType;
8485

@@ -305,6 +306,17 @@ export interface InputArrayType extends InputTypeBase {
305306
crossLanguageDefinitionId: string;
306307
}
307308

309+
export interface InputStreamingType extends InputTypeBase {
310+
kind: "streaming";
311+
name: string;
312+
valueType: InputType;
313+
streamKind: "jsonl" | "sse";
314+
contentTypes: string[];
315+
terminalEventType?: string;
316+
terminalEventValue?: string;
317+
crossLanguageDefinitionId: string;
318+
}
319+
308320
export interface InputDictionaryType extends InputTypeBase {
309321
kind: "dict";
310322
keyType: InputType;

packages/http-client-csharp/emitter/test/Unit/operation-converter.test.ts

Lines changed: 89 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -215,6 +215,95 @@ describe("Operation Converter", () => {
215215
});
216216

217217
describe("Operation response type conversion", () => {
218+
it("preserves JSONL stream item types", async () => {
219+
const program = await typeSpecCompile(
220+
`
221+
model Info {
222+
desc: string;
223+
}
224+
225+
@post op send(stream: JsonlStream<Info>): void;
226+
op receive(): JsonlStream<Info>;
227+
`,
228+
runner,
229+
);
230+
const context = createEmitterContext(program);
231+
const sdkContext = await createCSharpSdkContext(context);
232+
const [root] = createModel(sdkContext);
233+
234+
const send = root.clients[0].methods.find((method) => method.name === "send");
235+
ok(send);
236+
const sendParameter = send.parameters.find((parameter) => parameter.name === "stream");
237+
ok(sendParameter);
238+
ok(sendParameter.type.kind === "streaming");
239+
ok(sendParameter.type.valueType.kind === "model");
240+
strictEqual(sendParameter.type.valueType.name, "Info");
241+
242+
const bodyParameter = send.operation.parameters.find(
243+
(parameter) => parameter.kind === "body",
244+
);
245+
ok(bodyParameter);
246+
ok(bodyParameter.type.kind === "streaming");
247+
ok(bodyParameter.type.valueType.kind === "model");
248+
strictEqual(bodyParameter.type.valueType.name, "Info");
249+
250+
const receive = root.clients[0].methods.find((method) => method.name === "receive");
251+
ok(receive);
252+
ok(receive.response.type?.kind === "streaming");
253+
ok(receive.response.type.valueType.kind === "model");
254+
strictEqual(receive.response.type.valueType.name, "Info");
255+
ok(receive.operation.responses[0].bodyType?.kind === "streaming");
256+
ok(receive.operation.responses[0].bodyType.valueType.kind === "model");
257+
strictEqual(receive.operation.responses[0].bodyType.valueType.name, "Info");
258+
strictEqual(receive.operation.bufferResponse, false);
259+
});
260+
261+
it("preserves SSE event unions and terminal events", async () => {
262+
const program = await typeSpecCompile(
263+
`
264+
model ResponseCreated {
265+
id: string;
266+
}
267+
268+
model ResponseDelta {
269+
delta: string;
270+
}
271+
272+
@events
273+
union ResponseEvents {
274+
@Events.contentType("application/json")
275+
responseCreated: ResponseCreated,
276+
277+
@Events.contentType("application/json")
278+
responseDelta: ResponseDelta,
279+
280+
@Events.contentType("text/plain")
281+
@terminalEvent
282+
"[DONE]",
283+
}
284+
285+
op receive(): SSEStream<ResponseEvents>;
286+
`,
287+
runner,
288+
{ IsSseNeeded: true },
289+
);
290+
const context = createEmitterContext(program);
291+
const sdkContext = await createCSharpSdkContext(context);
292+
const [root] = createModel(sdkContext);
293+
294+
const receive = root.clients[0].methods.find((method) => method.name === "receive");
295+
ok(receive);
296+
ok(receive.response.type?.kind === "streaming");
297+
strictEqual(receive.response.type.streamKind, "sse");
298+
ok(receive.response.type.valueType.kind === "union");
299+
strictEqual(receive.response.type.valueType.name, "ResponseEvents");
300+
strictEqual(receive.response.type.terminalEventType, undefined);
301+
strictEqual(receive.response.type.terminalEventValue, "[DONE]");
302+
ok(receive.operation.responses[0].bodyType?.kind === "streaming");
303+
strictEqual(receive.operation.responses[0].bodyType.streamKind, "sse");
304+
strictEqual(receive.operation.bufferResponse, false);
305+
});
306+
218307
describe("With anonymous union enum response type", () => {
219308
it("should convert anonymous union enum response type to value type", async () => {
220309
const program = await typeSpecCompile(

packages/http-client-csharp/emitter/test/Unit/utils/test-util.ts

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,8 +4,11 @@ import { SdkTestLibrary } from "@azure-tools/typespec-client-generator-core/test
44
import type { CompilerOptions, EmitContext, Program } from "@typespec/compiler";
55
import type { TestHost } from "@typespec/compiler/testing";
66
import { createTestHost } from "@typespec/compiler/testing";
7+
import { EventsTestLibrary } from "@typespec/events/testing";
78
import { HttpTestLibrary } from "@typespec/http/testing";
89
import { RestTestLibrary } from "@typespec/rest/testing";
10+
import { SSETestLibrary } from "@typespec/sse/testing";
11+
import { StreamsTestLibrary } from "@typespec/streams/testing";
912
import { VersioningTestLibrary } from "@typespec/versioning/testing";
1013
import { XmlTestLibrary } from "@typespec/xml/testing";
1114
import { LoggerLevel } from "../../../src/lib/logger-level.js";
@@ -18,6 +21,9 @@ export async function createEmitterTestHost(): Promise<TestHost> {
1821
libraries: [
1922
RestTestLibrary,
2023
HttpTestLibrary,
24+
EventsTestLibrary,
25+
SSETestLibrary,
26+
StreamsTestLibrary,
2127
VersioningTestLibrary,
2228
AzureCoreTestLibrary,
2329
SdkTestLibrary,
@@ -45,6 +51,7 @@ export interface TypeSpecCompileOptions {
4551
AuthDecorator?: string;
4652
NoEmit?: boolean;
4753
IsVersionNeeded?: boolean;
54+
IsSseNeeded?: boolean;
4855
}
4956

5057
export async function typeSpecCompile(
@@ -57,6 +64,7 @@ export async function typeSpecCompile(
5764
const needTCGC = options?.IsTCGCNeeded ?? false;
5865
const needXml = options?.IsXmlNeeded ?? false;
5966
const needVersion = options?.IsVersionNeeded ?? true;
67+
const needSse = options?.IsSseNeeded ?? false;
6068
const authDecorator =
6169
options?.AuthDecorator ?? `@useAuth(ApiKeyAuth<ApiKeyLocation.header, "api-key">)`;
6270
const versions = `enum Versions {
@@ -77,12 +85,16 @@ export async function typeSpecCompile(
7785
const fileContent = `
7886
import "@typespec/rest";
7987
import "@typespec/http";
88+
import "@typespec/http/streams";
89+
${needSse ? 'import "@typespec/events";\nimport "@typespec/sse";' : ""}
8090
import "@typespec/versioning";
8191
${needXml ? 'import "@typespec/xml";' : ""}
8292
${needAzureCore ? 'import "@azure-tools/typespec-azure-core";' : ""}
8393
${needTCGC ? 'import "@azure-tools/typespec-client-generator-core";' : ""}
8494
using TypeSpec.Rest;
8595
using TypeSpec.Http;
96+
using TypeSpec.Http.Streams;
97+
${needSse ? "using TypeSpec.Events;\nusing TypeSpec.SSE;" : ""}
8698
using TypeSpec.Versioning;
8799
${needXml ? "using TypeSpec.Xml;" : ""}
88100
${needAzureCore ? "using Azure.Core;\nusing Azure.Core.Traits;" : ""}

packages/http-client-csharp/eng/scripts/Spector-Helper.psm1

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,6 @@ $failingSpecs = @(
55
Join-Path 'http' 'type' 'model' 'templated'
66
Join-Path 'http' 'type' 'file'
77
Join-Path 'http' 'client' 'naming' # pending until https://github.com/microsoft/typespec/issues/5653 is resolved
8-
Join-Path 'http' 'streaming' 'jsonl'
98
Join-Path 'http' 'type' 'union' 'discriminated' # pending design
109
Join-Path 'http' 'authentication' 'noauth' 'union' # pending design
1110
)

0 commit comments

Comments
 (0)