From 115d52eb943d2741c455a8c0078d3b4638981700 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Mike=20Kr=C3=BCger?= Date: Thu, 8 Oct 2026 11:12:11 +0200 Subject: [PATCH 1/2] Add bounded concurrent writes to import Default to 16 in-flight writes with configurable concurrency and sequential compatibility. Drain started writes on errors and cancellation, preserve accounting, and update help and documentation. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- CHANGELOG.md | 4 + .../CommandTests/ImportConcurrencyTests.cs | 371 ++++++++++++++++++ CosmosDBShell.Tests/ToolOperationsTests.cs | 12 + .../ImportCommand.cs | 192 ++++++--- CosmosDBShell/lang/en.ftl | 5 +- README.md | 2 +- docs/commands.md | 8 +- docs/mcp.md | 2 + l10n/CosmosDBShell.json | 5 +- 9 files changed, 544 insertions(+), 57 deletions(-) create mode 100644 CosmosDBShell.Tests/CommandTests/ImportConcurrencyTests.cs diff --git a/CHANGELOG.md b/CHANGELOG.md index 3034ec59..32904285 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,10 @@ ## Unreleased +### Improvements + +- `import` now streams documents with up to 16 concurrent writes by default. Use `--concurrency` to tune the limit or `--concurrency=1` for sequential, file-order writes. On a failure, imports stop scheduling new writes and wait for in-flight writes; `--continue-on-error` continues after per-item write failures. + ### Fixes - `for` and `do` loops now reject misspelled `in` and `while` keywords before any statements in the input execute. They previously accepted any identifier in those positions. diff --git a/CosmosDBShell.Tests/CommandTests/ImportConcurrencyTests.cs b/CosmosDBShell.Tests/CommandTests/ImportConcurrencyTests.cs new file mode 100644 index 00000000..13e6b347 --- /dev/null +++ b/CosmosDBShell.Tests/CommandTests/ImportConcurrencyTests.cs @@ -0,0 +1,371 @@ +// ------------------------------------------------------------ +// Copyright (c) Microsoft Corporation. All rights reserved. +// ------------------------------------------------------------ + +namespace CosmosShell.Tests.CommandTests; + +using System.Net; +using System.Text.Json; +using Azure.Data.Cosmos.Shell.Commands; +using Azure.Data.Cosmos.Shell.Core; +using Azure.Data.Cosmos.Shell.Parser; +using Microsoft.Azure.Cosmos; +using NSubstitute; + +public class ImportConcurrencyTests +{ + [Theory] + [InlineData(1)] + [InlineData(3)] + [InlineData(16)] + public async Task Writes_OverlapUpToLimitAndDrainBeforeReturning(int concurrency) + { + using var timeout = CreateTimeout(); + var container = Substitute.For(); + var gates = Enumerable.Range(0, concurrency + 2).Select(_ => NewGate()).ToArray(); + var started = Enumerable.Range(0, gates.Length).Select(_ => NewGate()).ToArray(); + var active = 0; + var peak = 0; + ConfigureWrites(container, async (item, token) => + { + var index = item.GetProperty("id").GetInt32(); + var current = Interlocked.Increment(ref active); + peak = Math.Max(peak, current); + started[index].TrySetResult(); + try + { + await gates[index].Task.WaitAsync(token); + return Response(HttpStatusCode.Created, 1.25); + } + finally + { + Interlocked.Decrement(ref active); + } + }); + + using var reader = NewReader(gates.Length); + var import = ImportCommand.WriteItemsAsync( + ImportCommand.EnumerateJsonLinesAsync(reader, timeout.Token), + container, ImportMode.Insert, concurrency, false, timeout.Token); + try + { + await Task.WhenAll(started.Take(concurrency).Select(gate => gate.Task)).WaitAsync(timeout.Token); + Assert.Equal(concurrency, active); + Assert.False(started[concurrency].Task.IsCompleted); + Assert.False(import.IsCompleted); + + // Completing the last write must free a slot even while earlier writes are blocked. + gates[concurrency - 1].TrySetResult(); + await started[concurrency].Task.WaitAsync(timeout.Token); + Assert.Equal(concurrency, active); + Assert.False(started[concurrency + 1].Task.IsCompleted); + } + finally + { + foreach (var gate in gates) + { + gate.TrySetResult(); + } + } + + var result = await import.WaitAsync(timeout.Token); + Assert.Equal((gates.Length, 0, gates.Length * 1.25), result); + Assert.Equal(concurrency, peak); + Assert.Equal(0, active); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task WriteFailure_StopsOrContinuesAndCountsInflightResults(bool continueOnError) + { + using var timeout = CreateTimeout(); + var container = Substitute.For(); + var first = new TaskCompletionSource>(TaskCreationOptions.RunContinuationsAsynchronously); + var second = new TaskCompletionSource>(TaskCreationOptions.RunContinuationsAsynchronously); + var started = new List(); + ConfigureWrites(container, (item, _) => + { + var id = item.GetProperty("id").GetInt32(); + started.Add(id); + if (id == 0) + { + return first.Task; + } + + if (id == 1) + { + return second.Task; + } + + if (id == 2) + { + return Task.FromException>( + new CosmosException("conflict", HttpStatusCode.Conflict, 0, "test", 3)); + } + + return Task.FromResult(Response(HttpStatusCode.Created, 2)); + }); + using var reader = NewReader(5); + var import = ImportCommand.WriteItemsAsync( + ImportCommand.EnumerateJsonLinesAsync(reader, timeout.Token), + container, ImportMode.Insert, 3, continueOnError, timeout.Token); + Assert.Equal(continueOnError ? 5 : 3, started.Count); + Assert.False(import.IsCompleted); + first.SetResult(Response(HttpStatusCode.Created, 2)); + second.SetResult(Response(HttpStatusCode.Created, 2)); + + var result = await import.WaitAsync(timeout.Token); + Assert.Equal(continueOnError ? (4, 1, 11.0) : (2, 1, 7.0), result); + Assert.Equal(continueOnError ? 5 : 3, started.Count); + } + + [Theory] + [InlineData(0, HttpStatusCode.Created, 1)] + [InlineData(0, HttpStatusCode.OK, 0)] + [InlineData(1, HttpStatusCode.Created, 1)] + [InlineData(1, HttpStatusCode.OK, 1)] + [InlineData(1, HttpStatusCode.BadRequest, 0)] + public async Task WriteMode_UsesCorrectApiAndCountsResponseStatus(int mode, HttpStatusCode status, int expectedSuccess) + { + using var timeout = CreateTimeout(); + var container = Substitute.For(); + var calls = 0; + ConfigureWrites(container, (_, _) => + { + calls++; + return Task.FromResult(Response(status, 2.5)); + }, (ImportMode)mode); + using var reader = NewReader(1); + var result = await ImportCommand.WriteItemsAsync( + ImportCommand.EnumerateJsonLinesAsync(reader, timeout.Token), + container, (ImportMode)mode, 16, false, timeout.Token); + + Assert.Equal(1, calls); + Assert.Equal((expectedSuccess, 1 - expectedSuccess, 2.5), result); + } + + [Fact] + public async Task SequentialWrites_PreserveFileOrderAndStopAtFailure() + { + using var timeout = CreateTimeout(); + var container = Substitute.For(); + var ids = new List(); + ConfigureWrites(container, (item, _) => + { + var id = item.GetProperty("id").GetInt32(); + ids.Add(id); + return Task.FromResult(Response(id == 2 ? HttpStatusCode.BadRequest : HttpStatusCode.Created, 1)); + }); + using var reader = NewReader(10); + var result = await ImportCommand.WriteItemsAsync( + ImportCommand.EnumerateJsonLinesAsync(reader, timeout.Token), + container, ImportMode.Insert, 1, false, timeout.Token); + + Assert.Equal(new[] { 0, 1, 2 }, ids); + Assert.Equal((2, 1, 3.0), result); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task ParseFailure_DrainsInflightWriteAndPreservesError(bool continueOnError) + { + using var timeout = CreateTimeout(); + var container = Substitute.For(); + var write = new TaskCompletionSource>(TaskCreationOptions.RunContinuationsAsynchronously); + var calls = 0; + ConfigureWrites(container, (_, _) => + { + calls++; + return write.Task; + }); + using var reader = new StringReader("{\"id\":0}\nnot-json\n{\"id\":2}\n"); + var import = ImportCommand.WriteItemsAsync( + ImportCommand.EnumerateJsonLinesAsync(reader, timeout.Token), + container, ImportMode.Insert, 3, continueOnError, timeout.Token); + Assert.Equal(1, calls); + Assert.False(import.IsCompleted); + write.SetResult(Response(HttpStatusCode.Created, 1)); + + var error = await Assert.ThrowsAsync(() => import.WaitAsync(timeout.Token)); + Assert.Contains("2", error.Message); + Assert.Equal(1, calls); + } + + [Fact] + public async Task Cancellation_PropagatesToWritesAndDrainsThem() + { + using var timeout = CreateTimeout(); + using var cancellation = CancellationTokenSource.CreateLinkedTokenSource(timeout.Token); + var container = Substitute.For(); + var bothStarted = NewGate(); + var active = 0; + var started = 0; + ConfigureWrites(container, async (_, token) => + { + Assert.Equal(cancellation.Token, token); + Interlocked.Increment(ref active); + if (Interlocked.Increment(ref started) == 2) + { + bothStarted.TrySetResult(); + } + + try + { + await Task.Delay(Timeout.InfiniteTimeSpan, token); + return Response(HttpStatusCode.Created, 1); + } + finally + { + Interlocked.Decrement(ref active); + } + }); + using var reader = NewReader(10); + var import = ImportCommand.WriteItemsAsync( + ImportCommand.EnumerateJsonLinesAsync(reader, cancellation.Token), + container, ImportMode.Insert, 2, true, cancellation.Token); + await bothStarted.Task.WaitAsync(timeout.Token); + await cancellation.CancelAsync(); + + await Assert.ThrowsAnyAsync(() => import.WaitAsync(timeout.Token)); + Assert.Equal(2, started); + Assert.Equal(0, active); + } + + [Fact] + public async Task UnexpectedWriteFailure_IsSurfacedAfterOtherWritesSettle() + { + using var timeout = CreateTimeout(); + var container = Substitute.For(); + var writes = Enumerable.Range(0, 2) + .Select(_ => new TaskCompletionSource>(TaskCreationOptions.RunContinuationsAsynchronously)) + .ToArray(); + ConfigureWrites(container, (item, _) => writes[item.GetProperty("id").GetInt32()].Task); + using var reader = NewReader(2); + var import = ImportCommand.WriteItemsAsync( + ImportCommand.EnumerateJsonLinesAsync(reader, timeout.Token), + container, ImportMode.Insert, 2, true, timeout.Token); + var expected = new InvalidOperationException("unexpected write failure"); + writes[0].SetException(expected); + Assert.False(import.IsCompleted); + writes[1].SetResult(Response(HttpStatusCode.Created, 1)); + + var actual = await Assert.ThrowsAsync(() => import.WaitAsync(timeout.Token)); + Assert.Same(expected, actual); + } + + [Fact] + public async Task Cancellation_WhileDrainingDoesNotReturnSuccessWhenWriteIgnoresToken() + { + using var timeout = CreateTimeout(); + using var cancellation = CancellationTokenSource.CreateLinkedTokenSource(timeout.Token); + var container = Substitute.For(); + var write = new TaskCompletionSource>(TaskCreationOptions.RunContinuationsAsynchronously); + ConfigureWrites(container, (_, _) => write.Task); + using var reader = NewReader(1); + var import = ImportCommand.WriteItemsAsync( + ImportCommand.EnumerateJsonLinesAsync(reader, cancellation.Token), + container, ImportMode.Insert, 2, false, cancellation.Token); + Assert.False(import.IsCompleted); + await cancellation.CancelAsync(); + Assert.False(import.IsCompleted); + write.SetResult(Response(HttpStatusCode.Created, 1)); + + await Assert.ThrowsAnyAsync(() => import.WaitAsync(timeout.Token)); + } + + [Theory] + [InlineData(0)] + [InlineData(-1)] + public async Task InvalidConcurrency_IsRejectedBeforeConnectionOrFileAccess(int concurrency) + { + var command = new ImportCommand { Concurrency = concurrency }; + var error = await Assert.ThrowsAsync(() => + command.ExecuteAsync(ShellInterpreter.Instance, new CommandState(), "import", TestContext.Current.CancellationToken)); + Assert.Contains("positive", error.Message, StringComparison.OrdinalIgnoreCase); + } + + [Fact] + public async Task EmptySource_DoesNotWrite() + { + using var reader = NewReader(0); + var container = Substitute.For(); + var result = await ImportCommand.WriteItemsAsync( + ImportCommand.EnumerateJsonLinesAsync(reader, TestContext.Current.CancellationToken), + container, ImportMode.Insert, 16, false, TestContext.Current.CancellationToken); + Assert.Equal((0, 0, 0.0), result); + Assert.Empty(container.ReceivedCalls()); + } + + [Theory] + [InlineData("jsonl", "{\"id\":\"a\"}\n{\"id\":\"b\"}\n")] + [InlineData("json", "[{\"id\":\"a\"},{\"id\":\"b\"}]")] + [InlineData("csv", "id\na\nb\n")] + public async Task DryRun_WithConcurrencyValidatesWithoutConnectionAndPreservesResult(string extension, string content) + { + var path = Path.Combine(Path.GetTempPath(), $"cosmos-import-{Guid.NewGuid():N}.{extension}"); + try + { + await File.WriteAllTextAsync(path, content, TestContext.Current.CancellationToken); + var command = new ImportCommand { File = path, DryRun = true, Concurrency = 32 }; + var state = await command.ExecuteAsync( + ShellInterpreter.Instance, new CommandState(), "import", TestContext.Current.CancellationToken); + var result = Assert.IsType(state.Result).Value; + + Assert.Equal("import", result.GetProperty("type").GetString()); + Assert.Equal(2, result.GetProperty("imported").GetInt32()); + Assert.Equal(0, result.GetProperty("failed").GetInt32()); + Assert.Equal(0, result.GetProperty("requestCharge").GetDouble()); + Assert.True(result.GetProperty("dryRun").GetBoolean()); + Assert.Null(state.RequestCharge); + } + finally + { + File.Delete(path); + } + } + + private static CancellationTokenSource CreateTimeout() + { + var timeout = CancellationTokenSource.CreateLinkedTokenSource(TestContext.Current.CancellationToken); + timeout.CancelAfter(TimeSpan.FromSeconds(10)); + return timeout; + } + + private static TaskCompletionSource NewGate() => new(TaskCreationOptions.RunContinuationsAsynchronously); + + private static StringReader NewReader(int count) => + new(string.Join('\n', Enumerable.Range(0, count).Select(id => JsonSerializer.Serialize(new { id })))); + + private static ItemResponse Response(HttpStatusCode status, double charge) + { + var response = Substitute.For>(); + response.StatusCode.Returns(status); + response.RequestCharge.Returns(charge); + return response; + } + + private static void ConfigureWrites( + Container container, + Func>> write, + ImportMode mode = ImportMode.Insert) + { + Task> Invoke(NSubstitute.Core.CallInfo call) + { + Assert.False(call.ArgAt(2).EnableContentResponseOnWrite); + return write(call.ArgAt(0), call.ArgAt(3)); + } + + if (mode == ImportMode.Upsert) + { + container.UpsertItemAsync(Arg.Any(), Arg.Any(), Arg.Any(), Arg.Any()) + .Returns(Invoke); + } + else + { + container.CreateItemAsync(Arg.Any(), Arg.Any(), Arg.Any(), Arg.Any()) + .Returns(Invoke); + } + } +} diff --git a/CosmosDBShell.Tests/ToolOperationsTests.cs b/CosmosDBShell.Tests/ToolOperationsTests.cs index 63a067f9..51b85d45 100644 --- a/CosmosDBShell.Tests/ToolOperationsTests.cs +++ b/CosmosDBShell.Tests/ToolOperationsTests.cs @@ -13,6 +13,18 @@ namespace CosmosShell.Tests; public class ToolOperationsTests { + [Fact] + public void GetTool_ImportExposesConfigurableConcurrencyWithDefault() + { + var factory = new CommandRunner().Commands["import"]; + var tool = ToolOperations.GetTool(factory); + var property = tool.InputSchema.GetProperty("properties").GetProperty("concurrency"); + + Assert.Equal("integer", property.GetProperty("type").GetString()); + Assert.Equal(16, property.GetProperty("default").GetInt32()); + Assert.Contains("sequential", property.GetProperty("description").GetString()); + } + [Fact] public void GetTool_IncludesCommandOptionsInInputSchema() { diff --git a/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/ImportCommand.cs b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/ImportCommand.cs index 0bf1a244..5bc7d5a7 100644 --- a/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/ImportCommand.cs +++ b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/ImportCommand.cs @@ -41,6 +41,7 @@ internal enum ImportMode [CosmosExample("import items.jsonl --mode=upsert", DescriptionKey = "command-import-example-5")] [CosmosExample("import items.jsonl --continue-on-error", DescriptionKey = "command-import-example-6")] [CosmosExample("import items.jsonl --dry-run", DescriptionKey = "command-import-example-7")] +[CosmosExample("import items.jsonl --concurrency=32", DescriptionKey = "command-import-example-8")] [McpAnnotation( Title = "Import Container Items", ReadOnly = false, @@ -49,6 +50,8 @@ internal enum ImportMode Description = "Bulk-loads items into a Cosmos container from a local JSON Lines, JSON array, or CSV file.")] internal class ImportCommand : CosmosCommand { + internal const int DefaultConcurrency = 16; + [CosmosParameter("file", RequiredErrorKey = "command-import-error-missing_file")] public string? File { get; init; } @@ -64,6 +67,9 @@ internal class ImportCommand : CosmosCommand [CosmosOption("format", "f", DefaultValue = ImportFormat.Auto)] public ImportFormat? Format { get; init; } + [CosmosOption("concurrency", DefaultValue = DefaultConcurrency)] + public int? Concurrency { get; init; } + [CosmosOption("partition-key", "pk")] public string? PartitionKey { get; init; } @@ -445,6 +451,12 @@ internal static JsonElement BuildCsvObject(IReadOnlyList headers, IReadO public override async Task ExecuteAsync(ShellInterpreter shell, CommandState commandState, string commandText, CancellationToken token) { + var concurrency = this.Concurrency ?? DefaultConcurrency; + if (concurrency < 1) + { + throw new CommandException("import", MessageService.GetString("command-import-error-invalid_concurrency")); + } + if (string.IsNullOrWhiteSpace(this.File)) { throw new CommandException("import", MessageService.GetString("command-import-error-missing_file")); @@ -483,7 +495,7 @@ public override async Task ExecuteAsync(ShellInterpreter shell, Co var continueOnError = this.ContinueOnError == true; var partitionKeySegments = ParsePartitionKeySegments(this.PartitionKey); - var (successCount, failCount, charge) = await ExecuteImportAsync(filePath, format, mode, container, continueOnError, dryRun, partitionKeySegments, token); + var (successCount, failCount, charge) = await ExecuteImportAsync(filePath, format, mode, container, continueOnError, dryRun, partitionKeySegments, concurrency, token); if (dryRun) { @@ -574,12 +586,9 @@ private static async Task ResolveFormatAsync(string filePath, Impo bool continueOnError, bool dryRun, string[]? partitionKeySegments, + int concurrency, CancellationToken token) { - var success = 0; - var failed = 0; - var charge = 0.0; - FileStream? stream = null; StreamReader? lineReader = null; IAsyncEnumerable<(int LineNumber, JsonElement Item)> source; @@ -603,75 +612,152 @@ private static async Task ResolveFormatAsync(string filePath, Impo } } - await foreach (var (lineNumber, item) in source.WithCancellation(token)) + if (dryRun) { - if (dryRun || container is null) + var count = 0; + await foreach (var item in source.WithCancellation(token)) { - success++; - continue; + count++; } - try + return (count, 0, 0); + } + + ArgumentNullException.ThrowIfNull(container); + return await WriteItemsAsync(source, container, mode, concurrency, continueOnError, token); + } + finally + { + lineReader?.Dispose(); + if (stream != null) + { + await stream.DisposeAsync(); + } + } + } + + internal static async Task<(int Success, int Failed, double Charge)> WriteItemsAsync( + IAsyncEnumerable<(int LineNumber, JsonElement Item)> source, + Container container, + ImportMode mode, + int concurrency, + bool continueOnError, + CancellationToken token) + { + ArgumentOutOfRangeException.ThrowIfLessThan(concurrency, 1); + var pending = new List>(); + var success = 0; + var failed = 0; + var charge = 0.0; + + void Record((double Charge, string? Error) result) + { + charge += result.Charge; + if (result.Error is null) + { + success++; + } + else + { + failed++; + ShellInterpreter.WriteLine(result.Error); + } + } + + async Task DrainCompletedAsync() + { + for (var i = pending.Count - 1; i >= 0; i--) + { + if (pending[i].IsCompleted) { - var response = mode == ImportMode.Upsert - ? await container.UpsertItemAsync(item, cancellationToken: token) - : await container.CreateItemAsync(item, cancellationToken: token); - charge += response.RequestCharge; + Record(await pending[i]); + pending.RemoveAt(i); + } + } + } - var ok = mode == ImportMode.Upsert - ? response.StatusCode == System.Net.HttpStatusCode.OK || response.StatusCode == System.Net.HttpStatusCode.Created - : response.StatusCode == System.Net.HttpStatusCode.Created; + await using var items = source.GetAsyncEnumerator(token); + try + { + while (true) + { + await DrainCompletedAsync(); + if (failed > 0 && !continueOnError) + { + break; + } - if (ok) - { - success++; - } - else - { - failed++; - ShellInterpreter.WriteLine(MessageService.GetArgsString( - "command-import-error-item_status", - "line", - lineNumber, - "status", - response.StatusCode.ToString())); - if (!continueOnError) - { - break; - } - } + if (pending.Count >= concurrency) + { + await Task.WhenAny(pending).WaitAsync(token); + continue; } - catch (CosmosException ce) + + token.ThrowIfCancellationRequested(); + if (!await items.MoveNextAsync()) { - failed++; - charge += ce.RequestCharge; - ShellInterpreter.WriteLine(MessageService.GetArgsString( - "command-import-error-item_failed", - "line", - lineNumber, - "status", - ce.StatusCode.ToString(), - "message", - CommandException.GetDisplayMessage(ce))); - if (!continueOnError) - { - break; - } + break; + } + + await DrainCompletedAsync(); + if (failed > 0 && !continueOnError) + { + break; } + + token.ThrowIfCancellationRequested(); + pending.Add(WriteItemAsync(container, mode, items.Current.LineNumber, items.Current.Item, token)); } } finally { - lineReader?.Dispose(); - if (stream != null) + // Drain every started write before returning or releasing the source, even on parse errors. + var results = await Task.WhenAll(pending); + foreach (var result in results) { - await stream.DisposeAsync(); + Record(result); } } + token.ThrowIfCancellationRequested(); return (success, failed, charge); } + private static async Task<(double Charge, string? Error)> WriteItemAsync( + Container container, + ImportMode mode, + int lineNumber, + JsonElement item, + CancellationToken token) + { + try + { + var options = new ItemRequestOptions { EnableContentResponseOnWrite = false }; + var response = mode == ImportMode.Upsert + ? await container.UpsertItemAsync(item, requestOptions: options, cancellationToken: token) + : await container.CreateItemAsync(item, requestOptions: options, cancellationToken: token); + var ok = response.StatusCode == System.Net.HttpStatusCode.Created || + (mode == ImportMode.Upsert && response.StatusCode == System.Net.HttpStatusCode.OK); + return (response.RequestCharge, ok ? null : MessageService.GetArgsString( + "command-import-error-item_status", + "line", + lineNumber, + "status", + response.StatusCode.ToString())); + } + catch (CosmosException ce) + { + return (ce.RequestCharge, MessageService.GetArgsString( + "command-import-error-item_failed", + "line", + lineNumber, + "status", + ce.StatusCode.ToString(), + "message", + CommandException.GetDisplayMessage(ce))); + } + } + private sealed class CancellationAwareTextReader(TextReader inner, CancellationToken token) : TextReader { public override int Peek() diff --git a/CosmosDBShell/lang/en.ftl b/CosmosDBShell/lang/en.ftl index d371835a..d1f85cc1 100644 --- a/CosmosDBShell/lang/en.ftl +++ b/CosmosDBShell/lang/en.ftl @@ -477,8 +477,9 @@ command-import-description-database = The database to write to. command-import-description-container = The container to write to. command-import-description-mode = Write mode: insert (default) or upsert. command-import-description-format = Input format: auto (default), jsonl, array, or csv. +command-import-description-concurrency = Maximum in-flight writes (default: 16). Must be positive. Use 1 for sequential writes and file-order processing. command-import-description-partition-key = For CSV import, the partition key path. Nested paths (e.g. /address/city) place the matching column under that path. -command-import-description-continue-on-error = Continue importing after individual item write failures. Parse or validation errors (invalid JSON, non-object rows, CSV partition-key conflicts) still abort the import. +command-import-description-continue-on-error = Continue importing after individual item write failures. Otherwise stop scheduling writes on the first observed failure and wait for in-flight writes. Parse or validation errors still abort the import. command-import-description-dry-run = Parse the file without writing any items. command-import-success = Imported { $count } { $count -> [one] item @@ -494,6 +495,7 @@ command-import-dry-run-success = Dry run: { $count } valid { $count -> *[other] items } command-import-error-missing_file = A source file path is required. +command-import-error-invalid_concurrency = Concurrency must be a positive integer. command-import-error-invalid_csv = Invalid CSV record at line { $line }. command-import-error-unnamed_csv_value = Line { $line }, column { $column } has a value but no CSV header. Add a column name before importing. script-error-argument-count = Function '{ $name }' expects { $expected } { $expected -> @@ -1624,6 +1626,7 @@ command-import-example-4 = Import CSV and nest the matching column under a neste command-import-example-5 = Insert new items and replace any existing items with the same id command-import-example-6 = Keep importing after individual item write failures command-import-example-7 = Validate the file without writing any items +command-import-example-8 = Import with up to 32 concurrent writes command-index-example-1 = Display the current container's indexing policy command-index-example-2 = Add a path to the included paths of the indexing policy command-index-example-3 = Remove a path from the indexing policy diff --git a/README.md b/README.md index c29368c3..7bdcff40 100644 --- a/README.md +++ b/README.md @@ -35,7 +35,7 @@ A terminal-native shell for Azure Cosmos DB — navigate databases like a filesy - MCP server for AI/tool integration - Distributed tracing via OpenTelemetry (`--otel`): emits a sampled W3C `traceparent` on Cosmos requests, with optional OTLP export -Exports replace their destination only after successful completion, preserving an existing file on failure or cancellation. Imports stream records; CSV exports use temporary disk storage to discover columns without retaining all documents in memory. Scalar CSV results use an empty column header by default; set `COSMOSDB_SHELL_CSV_SCALAR_COLUMN` to name it. An import rejects populated columns with empty headers rather than discarding their values. See [import/export](docs/commands.md#export). +Exports replace their destination only after successful completion, preserving an existing file on failure or cancellation. Imports stream records with up to 16 concurrent writes by default; set `--concurrency` to tune this, or use `--concurrency=1` for sequential, file-order writes. Concurrent writes can finish out of order, and in-flight writes may finish after an error. CSV exports use temporary disk storage to discover columns without retaining all documents in memory. Scalar CSV results use an empty column header by default; set `COSMOSDB_SHELL_CSV_SCALAR_COLUMN` to name it. An import rejects populated columns with empty headers rather than discarding their values. See [import/export](docs/commands.md#export). MCP command execution is serialized with the shell, and destructive confirmations are invalidated by connection or navigation changes. Ordinary explicit null MCP arguments are omitted; null continuation tokens and null `rm` partition-key/ETag safety options are rejected. MCP invocations are echoed in the shell so their activity stays visible, and they are recorded in history alongside interactive commands. Concurrent shells merge history under a shared lock and publish complete replacements instead of truncating the saved file. History remains fully replayable, including connection strings; treat its file as sensitive. See [MCP security](docs/mcp.md#security) and [history](docs/navigation.md#history). diff --git a/docs/commands.md b/docs/commands.md index 1c4c6f13..d14a523c 100644 --- a/docs/commands.md +++ b/docs/commands.md @@ -727,6 +727,8 @@ Options: --con, --container Target container (defaults to the current navigation context). --mode Write mode: insert (default) or upsert. --format, -f Input format: auto (default), jsonl, array, or csv. + --concurrency Maximum in-flight writes (default: 16; positive integer). + Use 1 for sequential writes in file order. --partition-key, --pk For CSV import, the partition key path. Nested paths (e.g. /address/city) nest the matching column. --continue, --continue-on-error @@ -737,13 +739,17 @@ Options: Examples: - `import items.jsonl` inserts every item from a JSON Lines file. +- `import items.jsonl --concurrency=32` allows up to 32 writes in flight. +- `import items.jsonl --concurrency=1` writes sequentially in file order. - `import items.json --format=array` reads a JSON array file. - `import items.csv` imports a CSV file, mapping each header column to a string property. - `import items.csv --partition-key=/address/city` nests the `city` column under `address` for a nested partition key. If a scalar column already occupies an intermediate path segment (for example an `address` column), the import fails with a conflict error rather than silently overwriting it. - `import items.jsonl --mode=upsert --continue-on-error` upserts items and keeps going on per-item failures. - `import items.jsonl --dry-run` validates the file without writing anything; useful before a real run. -By default, the first failure stops the import. With `--continue-on-error` the command keeps going after per-item *write* failures (for example a Cosmos write that throws) and the final summary reports how many items succeeded and how many failed. Parse and validation errors (invalid JSON, non-object rows, CSV partition-key conflicts) still abort the import immediately. The command exits with an error status if any items failed. +Imports keep up to 16 writes in flight by default without loading the entire file. `--concurrency` controls this bound; it must be a positive integer. Completion order is not guaranteed when concurrency exceeds 1. Use `--concurrency=1` when repeated IDs must be processed in file order (especially with upsert), or when exact sequential stop-on-error behavior is required. Higher concurrency can hide network latency but does not increase provisioned RU/s or avoid throttling; Cosmos SDK retry handling remains unchanged. Dry runs only validate input and do not start writes. + +By default, the first observed write failure stops scheduling new writes. Already-started writes are awaited and included in the final success/failure counts and request charge; they may succeed after the failure. With `--continue-on-error` the command keeps going after per-item *write* failures (for example a Cosmos write that throws) and the final summary reports how many items succeeded and how many failed. Parse and validation errors (invalid JSON, non-object rows, CSV partition-key conflicts) always stop scheduling writes, even with `--continue-on-error`, and in-flight writes are awaited before the error is returned. Cancellation is passed to both input reads and writes, and the command waits for started writes to settle before returning. Earlier writes are never rolled back. The command exits with an error status if any items failed. ### watch diff --git a/docs/mcp.md b/docs/mcp.md index a5c7d9fc..19304d01 100644 --- a/docs/mcp.md +++ b/docs/mcp.md @@ -97,6 +97,8 @@ On serverless accounts, `mkdb`, `mkcon`, and their `create` aliases omit through For deterministic ARM routing in multi-subscription environments, start the shell with `--connect-subscription` and `--connect-resource-group`. +MCP `import` uses the same bounded concurrency as the CLI: up to 16 writes in flight by default, configurable with the `concurrency` argument. Use `concurrency: 1` for sequential, file-order writes. In-flight writes may complete after an error; imports are not transactional. See [import](commands.md#import) for error and cancellation behavior. + ### Shell Location Updates Clients can read the `cosmos://shell/current-location` MCP resource. Its JSON content has a `currentLocation` field (`null` when disconnected, `/` at the account root, or `/database[/container]`) and a separate `currentAccountEndpoint` field (the connected Cosmos DB account URL, or `null` when disconnected). For example: `{"currentLocation":"/myDb/myContainer","currentAccountEndpoint":"https://myaccount.documents.azure.com/"}`. Clients that support resource subscriptions receive `notifications/resources/updated` when the shared shell location or connection changes, including changes made interactively. On notification, read the resource again for the new values; the notification itself contains only the URI. Rapid consecutive changes may be coalesced into a single notification. diff --git a/l10n/CosmosDBShell.json b/l10n/CosmosDBShell.json index 38a5fc5e..b119dc67 100644 --- a/l10n/CosmosDBShell.json +++ b/l10n/CosmosDBShell.json @@ -403,8 +403,9 @@ "command-help-example-3": "Show detailed help for all commands", "command-import-all-failed": "Failed to import all {0} items", "command-import-description": "Imports items into a container from a JSON Lines, JSON array, or CSV file.", + "command-import-description-concurrency": "Maximum in-flight writes (default: 16). Must be positive. Use 1 for sequential writes and file-order processing.", "command-import-description-container": "The container to write to.", - "command-import-description-continue-on-error": "Continue importing after individual item write failures. Parse or validation errors (invalid JSON, non-object rows, CSV partition-key conflicts) still abort the import.", + "command-import-description-continue-on-error": "Continue importing after individual item write failures. Otherwise stop scheduling writes on the first observed failure and wait for in-flight writes. Parse or validation errors still abort the import.", "command-import-description-database": "The database to write to.", "command-import-description-dry-run": "Parse the file without writing any items.", "command-import-description-file": "Source file path.", @@ -417,6 +418,7 @@ "command-import-error-blank_line": "Line {0} is blank.", "command-import-error-csv_pk_conflict": "CSV column \u0027{0}\u0027 conflicts with the partition key path \u0027{1}\u0027: the column holds a scalar value but the path requires it to be a nested object. Rename the column or choose a different partition key path.", "command-import-error-file_not_found": "File \u0027{0}\u0027 was not found.", + "command-import-error-invalid_concurrency": "Concurrency must be a positive integer.", "command-import-error-invalid_csv": "Invalid CSV record at line {0}.", "command-import-error-invalid_line_json": "Line {0} is not valid JSON: {1}", "command-import-error-item_failed": "Line {0}: {1} - {2}", @@ -432,6 +434,7 @@ "command-import-example-5": "Insert new items and replace any existing items with the same id", "command-import-example-6": "Keep importing after individual item write failures", "command-import-example-7": "Validate the file without writing any items", + "command-import-example-8": "Import with up to 32 concurrent writes", "command-import-success": "Imported {0} {1} (RU charge: {2})", "command-import-success-partial": "Imported {0} {1}, {2} failed (RU charge: {3})", "command-import-success-partial.__p1.one": "item", From 18cd3d9bc398bcccc1ae2e8f9cbd1c9ad137291c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Mike=20Kr=C3=BCger?= Date: Thu, 8 Oct 2026 11:42:47 +0200 Subject: [PATCH 2/2] Address concurrent import review feedback Expose reusable option minimum metadata in MCP schemas and keep dry-run test paths under the temporary directory. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- CosmosDBShell.Tests/CommandTests/ImportConcurrencyTests.cs | 3 ++- CosmosDBShell.Tests/ToolOperationsTests.cs | 2 ++ .../Azure.Data.Cosmos.Shell.Commands/ImportCommand.cs | 2 +- CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/Option.cs | 2 ++ .../Azure.Data.Cosmos.Shell.Core/CosmosOptionAttribute.cs | 2 ++ CosmosDBShell/Azure.Data.Cosmos.Shell.Mcp/ToolOperations.cs | 5 +++++ docs/mcp.md | 2 +- 7 files changed, 15 insertions(+), 3 deletions(-) diff --git a/CosmosDBShell.Tests/CommandTests/ImportConcurrencyTests.cs b/CosmosDBShell.Tests/CommandTests/ImportConcurrencyTests.cs index 13e6b347..8ebfa553 100644 --- a/CosmosDBShell.Tests/CommandTests/ImportConcurrencyTests.cs +++ b/CosmosDBShell.Tests/CommandTests/ImportConcurrencyTests.cs @@ -304,7 +304,8 @@ public async Task EmptySource_DoesNotWrite() [InlineData("csv", "id\na\nb\n")] public async Task DryRun_WithConcurrencyValidatesWithoutConnectionAndPreservesResult(string extension, string content) { - var path = Path.Combine(Path.GetTempPath(), $"cosmos-import-{Guid.NewGuid():N}.{extension}"); + var fileName = $"cosmos-import-{Guid.NewGuid():N}.{Path.GetFileName(extension).TrimStart('.')}"; + var path = Path.Join(Path.GetTempPath(), fileName); try { await File.WriteAllTextAsync(path, content, TestContext.Current.CancellationToken); diff --git a/CosmosDBShell.Tests/ToolOperationsTests.cs b/CosmosDBShell.Tests/ToolOperationsTests.cs index 51b85d45..4372a15c 100644 --- a/CosmosDBShell.Tests/ToolOperationsTests.cs +++ b/CosmosDBShell.Tests/ToolOperationsTests.cs @@ -22,6 +22,7 @@ public void GetTool_ImportExposesConfigurableConcurrencyWithDefault() Assert.Equal("integer", property.GetProperty("type").GetString()); Assert.Equal(16, property.GetProperty("default").GetInt32()); + Assert.Equal(1, property.GetProperty("minimum").GetInt32()); Assert.Contains("sequential", property.GetProperty("description").GetString()); } @@ -41,6 +42,7 @@ public void GetTool_IncludesCommandOptionsInInputSchema() Assert.Equal("string", queryProperty.GetProperty("type").GetString()); Assert.Equal("string", databaseProperty.GetProperty("type").GetString()); Assert.Equal("string", containerProperty.GetProperty("type").GetString()); + Assert.False(containerProperty.TryGetProperty("minimum", out _)); Assert.Equal("integer", maxProperty.GetProperty("type").GetString()); Assert.Equal(ToolOperations.DefaultPageSize, maxProperty.GetProperty("default").GetInt32()); Assert.Equal(1, maxProperty.GetProperty("minimum").GetInt32()); diff --git a/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/ImportCommand.cs b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/ImportCommand.cs index 5bc7d5a7..4d1ace93 100644 --- a/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/ImportCommand.cs +++ b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/ImportCommand.cs @@ -67,7 +67,7 @@ internal class ImportCommand : CosmosCommand [CosmosOption("format", "f", DefaultValue = ImportFormat.Auto)] public ImportFormat? Format { get; init; } - [CosmosOption("concurrency", DefaultValue = DefaultConcurrency)] + [CosmosOption("concurrency", DefaultValue = DefaultConcurrency, MinimumValue = 1)] public int? Concurrency { get; init; } [CosmosOption("partition-key", "pk")] diff --git a/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/Option.cs b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/Option.cs index 741ecb79..9f5d64b7 100644 --- a/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/Option.cs +++ b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/Option.cs @@ -28,6 +28,8 @@ public bool IsBool public object? DefaultValue => this.opt.DefaultValue; + public object? MinimumValue => this.opt.MinimumValue; + /// /// Gets the description of the option. /// diff --git a/CosmosDBShell/Azure.Data.Cosmos.Shell.Core/CosmosOptionAttribute.cs b/CosmosDBShell/Azure.Data.Cosmos.Shell.Core/CosmosOptionAttribute.cs index 3c564cd7..df5581ba 100644 --- a/CosmosDBShell/Azure.Data.Cosmos.Shell.Core/CosmosOptionAttribute.cs +++ b/CosmosDBShell/Azure.Data.Cosmos.Shell.Core/CosmosOptionAttribute.cs @@ -19,5 +19,7 @@ public CosmosOptionAttribute(params string[] name) public object? DefaultValue { get; set; } + public object? MinimumValue { get; set; } + public bool Hidden { get; set; } } diff --git a/CosmosDBShell/Azure.Data.Cosmos.Shell.Mcp/ToolOperations.cs b/CosmosDBShell/Azure.Data.Cosmos.Shell.Mcp/ToolOperations.cs index 5389c86a..4a74f554 100644 --- a/CosmosDBShell/Azure.Data.Cosmos.Shell.Mcp/ToolOperations.cs +++ b/CosmosDBShell/Azure.Data.Cosmos.Shell.Mcp/ToolOperations.cs @@ -136,6 +136,11 @@ internal static Tool GetTool(CommandFactory command) GetMcpOptionDescription(command, option), option.Name, GetMcpDefaultValue(command, option)); + if (option.MinimumValue is not null) + { + propertySchema["minimum"] = JsonSerializer.SerializeToNode(option.MinimumValue); + } + if (IsPagedMaxOption(command, option)) { propertySchema["minimum"] = 1; diff --git a/docs/mcp.md b/docs/mcp.md index 19304d01..9c4687af 100644 --- a/docs/mcp.md +++ b/docs/mcp.md @@ -97,7 +97,7 @@ On serverless accounts, `mkdb`, `mkcon`, and their `create` aliases omit through For deterministic ARM routing in multi-subscription environments, start the shell with `--connect-subscription` and `--connect-resource-group`. -MCP `import` uses the same bounded concurrency as the CLI: up to 16 writes in flight by default, configurable with the `concurrency` argument. Use `concurrency: 1` for sequential, file-order writes. In-flight writes may complete after an error; imports are not transactional. See [import](commands.md#import) for error and cancellation behavior. +MCP `import` uses the same bounded concurrency as the CLI: up to 16 writes in flight by default, configurable with the `concurrency` argument. Its schema requires a positive integer (`minimum: 1`). Use `concurrency: 1` for sequential, file-order writes. In-flight writes may complete after an error; imports are not transactional. See [import](commands.md#import) for error and cancellation behavior. ### Shell Location Updates