From f0a9b79be798673ef983f85d71672cf92c7fb9b2 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Mike=20Kr=C3=BCger?= Date: Fri, 9 Oct 2026 12:55:51 +0200 Subject: [PATCH] Add work-in-progress non-transactional bulk operations Add bounded concurrent writes, selection plans, journaling, interactive bulk state, and MCP integration. Related to #107. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- CHANGELOG.md | 4 + .../CommandTests/BulkCommandTests.cs | 863 ++++++++++++++++ CosmosDBShell.Tests/McpConfirmationTests.cs | 54 + .../Shell/CosmosShellPromptTests.cs | 13 + .../BatchOperationParser.cs | 4 +- .../BulkCommand.cs | 949 ++++++++++++++++++ .../BulkExecutor.cs | 141 +++ .../BulkJournal.cs | 165 +++ .../BulkOperation.cs | 30 + .../BulkOperationNormalizer.cs | 322 ++++++ .../BulkOutcome.cs | 44 + .../BulkSpool.cs | 88 ++ .../BulkSummary.cs | 57 ++ .../HelpCommand.cs | 4 +- .../CosmosShellPrompt.cs | 9 +- .../PendingBulkState.cs | 43 + .../ShellInterpreter.cs | 4 + .../ServerInstructions.md | 2 + .../ToolOperations.cs | 31 +- CosmosDBShell/lang/en.ftl | 86 ++ README.md | 1 + docs/commands.md | 123 ++- docs/mcp.md | 10 +- docs/navigation.md | 2 + l10n/CosmosDBShell.json | 81 ++ 25 files changed, 3122 insertions(+), 8 deletions(-) create mode 100644 CosmosDBShell.Tests/CommandTests/BulkCommandTests.cs create mode 100644 CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkCommand.cs create mode 100644 CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkExecutor.cs create mode 100644 CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkJournal.cs create mode 100644 CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkOperation.cs create mode 100644 CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkOperationNormalizer.cs create mode 100644 CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkOutcome.cs create mode 100644 CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkSpool.cs create mode 100644 CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkSummary.cs create mode 100644 CosmosDBShell/Azure.Data.Cosmos.Shell.Core/PendingBulkState.cs diff --git a/CHANGELOG.md b/CHANGELOG.md index ae499519..c1175eb6 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,10 @@ ## Unreleased +### New features + +- **Non-transactional `bulk` command.** Run unlimited independent `create`, `upsert`, `replace`, `delete`, and `patch` operations across partition keys with bounded concurrency (default 16). It uses the batch operation schema, plus optional per-operation `partitionKey` and `ifMatch`, and the same `run`/`begin`/`add`/`execute`/`cancel`/`status`/`show` subcommands. `bulk patch --where` and `bulk delete --where` select items by predicate, including complete hierarchical partition keys. Add `--save` to write a reviewable plan for `bulk run`. Safeguards include dry-run, confirmation, item and observed-RU limits, and ETag checks. Journals skip succeeded writes on rerun, retry failed writes, and hold back writes with unknown outcomes. MCP exposes `run`, `patch`, and `delete`. ([#107](https://github.com/Azure/CosmosDBShell/issues/107)) + ### Fixes - Parse `exec` options and shell words like direct commands. Bind built-in options normally and pass option-shaped words to functions and script files as positional text. diff --git a/CosmosDBShell.Tests/CommandTests/BulkCommandTests.cs b/CosmosDBShell.Tests/CommandTests/BulkCommandTests.cs new file mode 100644 index 00000000..a275b369 --- /dev/null +++ b/CosmosDBShell.Tests/CommandTests/BulkCommandTests.cs @@ -0,0 +1,863 @@ +// ------------------------------------------------------------ +// Copyright (c) Microsoft Corporation. All rights reserved. +// ------------------------------------------------------------ + +namespace CosmosShell.Tests.CommandTests; + +using System.Globalization; +using System.Net; +using System.Runtime.CompilerServices; +using System.Text; +using System.Text.Json; +using Azure.Data.Cosmos.Shell.Commands; +using Azure.Data.Cosmos.Shell.Core; +using Azure.Data.Cosmos.Shell.Mcp; +using Azure.Data.Cosmos.Shell.Parser; +using Azure.Data.Cosmos.Shell.States; +using Microsoft.Azure.Cosmos; +using NSubstitute; + +public class BulkCommandTests +{ + private static readonly string[] SinglePath = ["/pk"]; + + private static readonly string[] HierarchicalPaths = ["/tenant/id", "/region"]; + + // Normalization: batch schema plus per-operation partition keys. + [Fact] + public void NormalizeAcceptsEveryBatchOperationKind() + { + string[] operations = + [ + """{"op":"create","item":{"id":"1","pk":"a"}}""", + """{"op":"upsert","item":{"id":"2","pk":"a"}}""", + """{"op":"replace","id":"3","item":{"id":"3","pk":"a"},"ifMatch":"tag"}""", + """{"op":"delete","id":"4","partitionKey":"a"}""", + """{"op":"patch","id":"5","partitionKey":"a","operations":[{"op":"set","path":"/v","value":2}]}""", + ]; + + var normalized = operations.Select((json, index) => BulkOperationNormalizer.Normalize(Json(json), SinglePath, null, index)).ToList(); + + Assert.Equal(["create", "upsert", "replace", "delete", "patch"], normalized.Select(operation => operation.Op)); + Assert.All(normalized, operation => Assert.Equal("""["a"]""", operation.PartitionKey)); + Assert.Contains("\"ifMatch\":\"tag\"", normalized[2].Json); + } + + [Fact] + public void NormalizeIsIdempotentSoSavedPlansHashIdentically() + { + var first = BulkOperationNormalizer.Normalize( + Json("""{ "operations": [ {"op":"incr","path":"/n","value":1.50} ], "partitionKey": ["t", 4], "id":"x", "op":"patch" }"""), + HierarchicalPaths, + null, + 0); + var second = BulkOperationNormalizer.Normalize(Json(first.Json), HierarchicalPaths, null, 0); + + Assert.Equal(first.Json, second.Json); + Assert.Equal(BulkSpool.ComputeHash([first.Json]), BulkSpool.ComputeHash([second.Json])); + Assert.Contains("1.50", first.Json); + } + + [Fact] + public void ItemOperationsDeriveHierarchicalKeysPreservingTypesAndUndefined() + { + var operation = BulkOperationNormalizer.Normalize(Json("""{"op":"upsert","item":{"id":"a","tenant":{"id":"42"}}}"""), HierarchicalPaths, null, 0); + Assert.Equal("""["42",{}]""", operation.PartitionKey); + Assert.Equal( + new PartitionKeyBuilder().Add("42").AddNoneType().Build(), + BulkOperationNormalizer.ParseKey(operation.PartitionKey)); + + var nullKey = BulkOperationNormalizer.Normalize(Json("""{"op":"upsert","item":{"id":"a","pk":null}}"""), SinglePath, null, 0); + Assert.Equal(PartitionKey.Null, BulkOperationNormalizer.ParseKey(nullKey.PartitionKey)); + var missingKey = BulkOperationNormalizer.Normalize(Json("""{"op":"upsert","item":{"id":"a"}}"""), SinglePath, null, 0); + Assert.Equal(PartitionKey.None, BulkOperationNormalizer.ParseKey(missingKey.PartitionKey)); + } + + [Fact] + public void DefaultPartitionKeyAppliesToAddressedOperationsAndIsCheckedForItems() + { + var defaultKey = BulkOperationNormalizer.CanonicalKeyFromArgument("tenant", 1); + var delete = BulkOperationNormalizer.Normalize(Json("""{"op":"delete","id":"1"}"""), SinglePath, defaultKey, 0); + Assert.Equal("""["tenant"]""", delete.PartitionKey); + + var explicitKey = BulkOperationNormalizer.Normalize(Json("""{"op":"delete","id":"1","partitionKey":"other"}"""), SinglePath, defaultKey, 0); + Assert.Equal("""["other"]""", explicitKey.PartitionKey); + + var error = Assert.Throws(() => + BulkOperationNormalizer.Normalize(Json("""{"op":"upsert","item":{"id":"1","pk":"other"}}"""), SinglePath, defaultKey, 6)); + Assert.StartsWith("Operation 7:", error.Message); + BulkOperationNormalizer.Normalize(Json("""{"op":"upsert","item":{"id":"1","pk":1}}"""), SinglePath, BulkOperationNormalizer.CanonicalKeyFromArgument("1.0", 1), 0); + } + + [Theory] + [InlineData("""{"op":"delete","id":"1"}""")] + [InlineData("""{"op":"delete","id":"1","partitionKey":["a"]}""")] + [InlineData("""{"op":"delete","id":"1","partitionKey":[{"x":1},"a"]}""")] + [InlineData("""{"op":"delete","id":1,"partitionKey":["a","b"]}""")] + [InlineData("""{"op":"create","item":{"id":"1"},"ifMatch":"tag"}""")] + [InlineData("""{"op":"replace","id":"2","item":{"id":"1"}}""")] + [InlineData("""{"op":"patch","id":"1","partitionKey":["a","b"],"operations":[{"op":"remove","path":"old"}]}""")] + [InlineData("""{"op":"patch","id":"1","partitionKey":["a","b"],"operations":[]}""")] + [InlineData("""{"op":"unknown","id":"1","partitionKey":["a","b"]}""")] + public void InvalidOperationsAreRejected(string json) + { + Assert.Throws(() => BulkOperationNormalizer.Normalize(Json(json), HierarchicalPaths, null, 0)); + } + + [Fact] + public void PatchSupportsAtMostTenOperations() + { + var ten = JsonSerializer.SerializeToElement(Enumerable.Repeat(new { op = "remove", path = "/old" }, 10)); + BulkOperationNormalizer.ValidatePatchOperations(ten); + var eleven = JsonSerializer.SerializeToElement(Enumerable.Repeat(new { op = "remove", path = "/old" }, 11)); + Assert.Throws(() => BulkOperationNormalizer.ValidatePatchOperations(eleven)); + } + + // Option validation mirrors the per-subcommand surface. + [Theory] + [InlineData("")] + [InlineData("unknown")] + public void InvalidSubcommandsAreRejected(string subcommand) + { + Assert.Throws(() => new BulkCommand { Subcommand = subcommand }.Validate(BulkCommand.NormalizeSubcommand(subcommand))); + } + + [Fact] + public void SubcommandAliasesMatchBatch() + { + Assert.Equal("execute", BulkCommand.NormalizeSubcommand("EXEC")); + Assert.Equal("execute", BulkCommand.NormalizeSubcommand("commit")); + Assert.Equal("cancel", BulkCommand.NormalizeSubcommand("abort")); + } + + public static TheoryData InvalidOptionCombinationIndexes => new(Enumerable.Range(0, InvalidOptionCombinations.Length)); + + private static (string Subcommand, BulkCommand Command)[] InvalidOptionCombinations => + [ + ("run", new BulkCommand { Where = "true" }), + ("run", new BulkCommand { ETag = true }), + ("run", new BulkCommand { Save = "plan.jsonl" }), + ("run", new BulkCommand { RetryUncertain = true }), + ("run", new BulkCommand { DryRun = true, Journal = "journal.jsonl" }), + ("patch", new BulkCommand { Where = "true" }), + ("patch", new BulkCommand { Where = "true", Operations = "not json" }), + ("patch", new BulkCommand { Where = "true", Operations = """[{"op":"set","path":"/v"}]""" }), + ("delete", new BulkCommand { Data = "SELECT * FROM c" }), + ("delete", new BulkCommand { Where = "true", Operations = "[]" }), + ("delete", new BulkCommand { Where = "true", Journal = "journal.jsonl" }), + ("delete", new BulkCommand { Where = " " }), + ("begin", new BulkCommand { Yes = true }), + ("execute", new BulkCommand { PartitionKeyArgument = "a" }), + ("status", new BulkCommand { Database = "db" }), + ("run", new BulkCommand { Concurrency = 0 }), + ("delete", new BulkCommand { Where = "true", MaxItems = 0 }), + ("run", new BulkCommand { MaxRu = double.NaN }), + ("run", new BulkCommand { MaxRu = -1 }), + ]; + + [Theory] + [MemberData(nameof(InvalidOptionCombinationIndexes))] + public void InvalidOptionCombinationsAreRejected(int index) + { + var (subcommand, command) = InvalidOptionCombinations[index]; + Assert.Throws(() => command.Validate(subcommand)); + } + + [Fact] + public void SelectionQueryProjectsEscapedPartitionKeyPaths() + { + var query = BulkCommand.BuildSelectionQuery("c.expired = true", ["/tenant/id", "/re\"gion"]); + Assert.Equal("""SELECT c.id, c._etag, c["tenant"]["id"] AS __pk0, c["re\"gion"] AS __pk1 FROM c WHERE (c.expired = true)""", query); + } + + // Writes. + [Theory] + [InlineData("create")] + [InlineData("upsert")] + [InlineData("replace")] + [InlineData("delete")] + [InlineData("patch")] + public async Task WriteUsesMatchingApiWithKeyIfMatchAndNoContentResponse(string op) + { + var container = Substitute.For(); + Fixture.ConfigureWrites(container, HttpStatusCode.OK, 3.5); + var json = op switch + { + "create" => """{"op":"create","item":{"id":"1","pk":"a"}}""", + "upsert" or "replace" => $$"""{"op":"{{op}}","item":{"id":"1","pk":"a"},"ifMatch":"tag"}""", + "delete" => """{"op":"delete","id":"1","partitionKey":"a","ifMatch":"tag"}""", + _ => """{"op":"patch","id":"1","partitionKey":"a","ifMatch":"tag","operations":[{"op":"set","path":"/v","value":2}]}""", + }; + var operation = BulkOperationNormalizer.Normalize(Json(json), SinglePath, null, 0); + + var result = await BulkCommand.WriteAsync(container, operation, TestContext.Current.CancellationToken); + + Assert.Equal(BulkOutcome.Succeeded, result.Status); + Assert.Equal(3.5, result.RequestCharge); + var call = Assert.Single(container.ReceivedCalls()); + Assert.StartsWith(char.ToUpperInvariant(op[0]) + op[1..], call.GetMethodInfo().Name); + var arguments = call.GetArguments(); + Assert.Contains(arguments, argument => argument is PartitionKey key && key.Equals(new PartitionKey("a"))); + var options = Assert.Single(arguments.OfType()); + Assert.False(options.EnableContentResponseOnWrite); + Assert.Equal(op == "create" ? null : "tag", options.IfMatchEtag); + } + + [Theory] + [InlineData(404, BulkOutcome.Failed)] + [InlineData(412, BulkOutcome.Failed)] + [InlineData(429, BulkOutcome.Failed)] + [InlineData(408, BulkOutcome.Uncertain)] + [InlineData(503, BulkOutcome.Uncertain)] + public async Task WriteClassifiesServiceErrors(int status, string expected) + { + var container = Substitute.For(); + container.DeleteItemAsync(default!, default, default, TestContext.Current.CancellationToken) + .ReturnsForAnyArgs>>(_ => throw new CosmosException("failure", (HttpStatusCode)status, 0, "activity", 2)); + + var result = await BulkCommand.WriteAsync(container, Operation(0), TestContext.Current.CancellationToken); + + Assert.Equal(expected, result.Status); + Assert.Equal(status, result.StatusCode); + Assert.Equal(2, result.RequestCharge); + } + + // Execution engine. + [Theory] + [InlineData(1)] + [InlineData(3)] + [InlineData(16)] + public async Task ConcurrencyRefillsWhicheverSlotFinishesAndDrains(int concurrency) + { + using var timeout = Timeout(); + var gates = Enumerable.Range(0, concurrency + 2).Select(_ => new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously)).ToArray(); + var started = Enumerable.Range(0, gates.Length).Select(_ => new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously)).ToArray(); + var summary = new BulkSummary(); + var run = BulkExecutor.ExecuteAsync( + SourceAsync(gates.Length, timeout.Token), + (operation, _) => + { + started[(int)operation.Index].SetResult(); + return gates[(int)operation.Index].Task; + }, + IgnoreAsync, + new Dictionary(), + summary, + concurrency, + null, + false, + false, + timeout.Token); + try + { + await Task.WhenAll(started.Take(concurrency).Select(gate => gate.Task)).WaitAsync(timeout.Token); + Assert.False(started[concurrency].Task.IsCompleted); + gates[concurrency - 1].SetResult(Succeeded(concurrency - 1)); + await started[concurrency].Task.WaitAsync(timeout.Token); + Assert.False(started[concurrency + 1].Task.IsCompleted); + Assert.False(run.IsCompleted); + } + finally + { + for (var i = 0; i < gates.Length; i++) + { + gates[i].TrySetResult(Succeeded(i)); + } + } + + await run.WaitAsync(timeout.Token); + Assert.Equal(gates.Length, summary.Attempted); + Assert.Equal(gates.Length, summary.Succeeded); + Assert.Equal(gates.Length, summary.RequestCharge); + Assert.True(summary.Success); + } + + [Theory] + [InlineData(false, 1)] + [InlineData(true, 5)] + public async Task FailureStopsOrContinues(bool continueOnError, int attempted) + { + var summary = new BulkSummary(); + await BulkExecutor.ExecuteAsync( + SourceAsync(5, TestContext.Current.CancellationToken), + (operation, _) => Task.FromResult(Succeeded(operation.Index) with { Status = BulkOutcome.Failed, StatusCode = 409 }), + IgnoreAsync, + new Dictionary(), + summary, + 1, + null, + continueOnError, + false, + TestContext.Current.CancellationToken); + + Assert.Equal(attempted, summary.Attempted); + Assert.Equal(attempted, summary.Failed); + Assert.Equal(!continueOnError, summary.ResultIncomplete); + Assert.False(summary.Success); + } + + [Fact] + public async Task RuBudgetStopsNewWrites() + { + var summary = new BulkSummary { RequestCharge = 2 }; + await BulkExecutor.ExecuteAsync( + SourceAsync(5, TestContext.Current.CancellationToken), + (operation, _) => Task.FromResult(Succeeded(operation.Index)), + IgnoreAsync, + new Dictionary(), + summary, + 1, + 4, + true, + false, + TestContext.Current.CancellationToken); + + Assert.Equal(2, summary.Attempted); + Assert.True(summary.BudgetExceeded); + Assert.True(summary.ResultIncomplete); + } + + [Fact] + public async Task OperationsOnTheSameItemRunInOrderEvenWithEquivalentKeys() + { + using var timeout = Timeout(); + var gate = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var started = new List(); + async IAsyncEnumerable OperationsAsync() + { + yield return new BulkOperation(0, "patch", "same", "[1]", "{}"); + yield return new BulkOperation(1, "patch", "same", "[1.0]", "{}"); + await Task.CompletedTask; + } + + var run = BulkExecutor.ExecuteAsync( + OperationsAsync(), + (operation, _) => + { + started.Add(operation.Index); + return operation.Index == 0 ? gate.Task : Task.FromResult(Succeeded(1)); + }, + IgnoreAsync, + new Dictionary(), + new BulkSummary(), + 16, + null, + true, + false, + timeout.Token); + + Assert.Equal([0L], started); + gate.SetResult(Succeeded(0)); + await run.WaitAsync(timeout.Token); + Assert.Equal([0L, 1L], started); + } + + [Fact] + public async Task CancellationDrainsInFlightWritesAndRecordsUncertainOutcome() + { + using var cancellation = Timeout(); + var started = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var records = new List(); + var summary = new BulkSummary(); + var run = BulkExecutor.ExecuteAsync( + SourceAsync(3, cancellation.Token), + async (operation, token) => + { + started.TrySetResult(); + await Task.Delay(System.Threading.Timeout.Infinite, token); + return Succeeded(operation.Index); + }, + outcome => + { + records.Add(outcome); + return Task.CompletedTask; + }, + new Dictionary(), + summary, + 1, + null, + true, + false, + cancellation.Token); + + await started.Task.WaitAsync(TestContext.Current.CancellationToken); + await cancellation.CancelAsync(); + await Assert.ThrowsAnyAsync(async () => await run.WaitAsync(TestContext.Current.CancellationToken)); + Assert.Equal(1, summary.Attempted); + Assert.Equal(1, summary.Uncertain); + Assert.Equal([BulkOutcome.Started, BulkOutcome.Uncertain], records.Select(record => record.Status)); + } + + [Fact] + public async Task UnexpectedFailureDrainsOtherWritesAndKeepsOriginalException() + { + using var timeout = Timeout(); + var first = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var second = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var summary = new BulkSummary(); + var run = BulkExecutor.ExecuteAsync( + SourceAsync(2, timeout.Token), + (operation, _) => operation.Index == 0 ? first.Task : second.Task, + IgnoreAsync, + new Dictionary(), + summary, + 2, + null, + true, + false, + timeout.Token); + var expected = new IOException("unexpected"); + first.SetException(expected); + Assert.False(run.IsCompleted); + second.SetResult(Succeeded(1)); + + var actual = await Assert.ThrowsAsync(async () => await run.WaitAsync(timeout.Token)); + Assert.Same(expected, actual); + Assert.Equal(1, summary.Succeeded); + Assert.Equal(1, summary.Uncertain); + } + + [Fact] + public async Task ResumeSkipsSucceededRetriesFailedAndHoldsBackUncertain() + { + var previous = new Dictionary + { + [0] = Succeeded(0), + [1] = Succeeded(1) with { Status = BulkOutcome.Failed, StatusCode = 429 }, + [2] = Succeeded(2) with { Status = BulkOutcome.Started }, + [3] = Succeeded(3) with { Status = BulkOutcome.Uncertain }, + }; + var attempted = new List(); + var summary = new BulkSummary(); + await BulkExecutor.ExecuteAsync( + SourceAsync(5, TestContext.Current.CancellationToken), + (operation, _) => + { + attempted.Add(operation.Index); + return Task.FromResult(Succeeded(operation.Index)); + }, + IgnoreAsync, + previous, + summary, + 1, + null, + false, + false, + TestContext.Current.CancellationToken); + + Assert.Equal([1L, 4L], attempted); + Assert.Equal(3, summary.Skipped); + Assert.Equal(2, summary.Failed); + Assert.Equal(2, summary.Uncertain); + Assert.False(summary.Success); + + attempted.Clear(); + await BulkExecutor.ExecuteAsync( + SourceAsync(5, TestContext.Current.CancellationToken), + (operation, _) => + { + attempted.Add(operation.Index); + return Task.FromResult(Succeeded(operation.Index)); + }, + IgnoreAsync, + previous, + new BulkSummary(), + 1, + null, + false, + true, + TestContext.Current.CancellationToken); + Assert.Equal([1L, 2L, 3L, 4L], attempted); + } + + // Journal. + [Fact] + public async Task JournalPersistsOutcomesTrimsTornTailAndRejectsMismatchOrConcurrentUse() + { + using var directory = new TempDirectory(); + var path = directory.PathFor("journal.jsonl"); + using (var journal = BulkJournal.Open(path, "endpoint", "rid", "hash", 3)) + { + await journal.RecordAsync(Succeeded(0) with { Status = BulkOutcome.Started }); + await journal.RecordAsync(Succeeded(0)); + await journal.RecordAsync(Succeeded(1) with { Status = BulkOutcome.Started }); + Assert.ThrowsAny(() => BulkJournal.Open(path, "endpoint", "rid", "hash", 3)); + } + + await File.AppendAllTextAsync(path, "{\"index\":2,\"op\":\"del", TestContext.Current.CancellationToken); + using (var journal = BulkJournal.Open(path, "endpoint", "rid", "hash", 3)) + { + Assert.Equal(BulkOutcome.Succeeded, journal.Previous[0].Status); + Assert.Equal(BulkOutcome.Started, journal.Previous[1].Status); + Assert.False(journal.Previous.ContainsKey(2)); + await journal.RecordAsync(Succeeded(2)); + } + + Assert.EndsWith("\n", await File.ReadAllTextAsync(path, TestContext.Current.CancellationToken)); + Assert.Throws(() => BulkJournal.Open(path, "endpoint", "rid", "other-hash", 3)); + Assert.Throws(() => BulkJournal.Open(path, "endpoint", "recreated", "hash", 3)); + + var lines = await File.ReadAllLinesAsync(path, TestContext.Current.CancellationToken); + lines[1] = "not json"; + await File.WriteAllTextAsync(path, string.Join("\n", lines) + "\n", TestContext.Current.CancellationToken); + Assert.Throws(() => BulkJournal.Open(path, "endpoint", "rid", "hash", 3)); + } + + // End-to-end command behavior against a substituted container. + [Fact] + public async Task RunExecutesMixedOperationsAcrossPartitions() + { + using var fixture = new Fixture(); + var state = await fixture.RunAsync( + """bulk run '[{"op":"upsert","item":{"id":"1","pk":"a"}},{"op":"delete","id":"2","partitionKey":"b"},{"op":"patch","id":"3","partitionKey":"c","operations":[{"op":"set","path":"/v","value":2}]}]' --yes --db db --con items"""); + + var result = Assert.IsType(state.Result).Value; + Assert.False(state.IsError); + Assert.Equal(3, result.GetProperty("operationCount").GetInt64()); + Assert.Equal(3, result.GetProperty("succeeded").GetInt64()); + Assert.True(result.GetProperty("success").GetBoolean()); + Assert.Equal(17, state.RequestCharge); + } + + [Fact] + public async Task RunWithPipedOperationsAndDefaultPartitionKeyMatchesBatchStyle() + { + using var fixture = new Fixture(); + var state = await fixture.RunAsync("""echo '[{"op":"delete","id":"1"},{"op":"delete","id":"2"}]' | bulk run --partition-key a --yes --db db --con items"""); + + Assert.False(state.IsError); + await fixture.Container.Received(2).DeleteItemAsync(Arg.Any(), new PartitionKey("a"), Arg.Any(), Arg.Any()); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task RunReadsFilesAndValidatesEverythingBeforeWriting(bool array) + { + using var fixture = new Fixture(); + var path = fixture.Directory.PathFor("operations.json"); + string[] rows = ["""{"op":"upsert","item":{"id":"1","pk":"a"}}""", """{"op":"delete","id":"2"}"""]; + await File.WriteAllTextAsync(path, array ? "[" + string.Join(",", rows) + "]" : string.Join("\n", rows), TestContext.Current.CancellationToken); + + var exception = await Assert.ThrowsAsync(() => fixture.RunAsync($"bulk run '{path}' --yes --db db --con items")); + + Assert.StartsWith("Operation 2:", exception.Message); + Assert.Empty(fixture.Writes()); + } + + [Fact] + public async Task NonInteractiveWritesRequireYesBeforeAnyRequest() + { + using var fixture = new Fixture(); + await Assert.ThrowsAsync(() => fixture.RunAsync("""bulk run '[{"op":"delete","id":"1","partitionKey":"a"}]' --db db --con items""")); + Assert.Empty(fixture.Container.ReceivedCalls()); + } + + [Fact] + public async Task DryRunValidatesWithoutWriting() + { + using var fixture = new Fixture(); + var state = await fixture.RunAsync("""bulk run '[{"op":"delete","id":"1","partitionKey":"a"}]' --dry-run --db db --con items"""); + + var result = Assert.IsType(state.Result).Value; + Assert.True(result.GetProperty("dryRun").GetBoolean()); + Assert.Equal(1, result.GetProperty("operationCount").GetInt64()); + Assert.Empty(fixture.Writes()); + } + + [Fact] + public async Task WhereSelectionSavesPlanThatRunResumesWithJournal() + { + using var fixture = new Fixture(); + fixture.Query( + """{"id":"a","_etag":"e1","__pk0":"tenant"}""", + """{"id":"b","_etag":"e2"}""", + """{"id":"c","_etag":"e3","__pk0":null}"""); + var plan = fixture.Directory.PathFor("plan.jsonl"); + var journal = fixture.Directory.PathFor("journal.jsonl"); + + var preview = await fixture.RunAsync( + $$"""bulk patch --where "c.v = 1" --operations '[{"op":"set","path":"/v","value":2}]' --etag --dry-run --save '{{plan}}' --db db --con items"""); + Assert.Equal(3, Assert.IsType(preview.Result).Value.GetProperty("operationCount").GetInt64()); + Assert.Contains("c.v = 1", fixture.LastQuery); + Assert.Empty(fixture.Writes()); + + var saved = await File.ReadAllLinesAsync(plan, TestContext.Current.CancellationToken); + Assert.Equal(3, saved.Length); + Assert.Contains("\"partitionKey\":[{}]", saved[1]); + Assert.Contains("\"partitionKey\":[null]", saved[2]); + Assert.Contains("\"ifMatch\":\"e1\"", saved[0]); + + fixture.FailPatchFor("b", HttpStatusCode.TooManyRequests); + var first = await fixture.RunAsync($"bulk run '{plan}' --journal '{journal}' --continue-on-error --yes --db db --con items"); + Assert.True(first.IsError); + Assert.Equal(ShellExitCode.GeneralFailure, first.ExitCode); + + fixture.FailPatchFor(null, HttpStatusCode.OK); + var resumed = await fixture.RunAsync($"bulk run '{plan}' --journal '{journal}' --yes --db db --con items"); + var result = Assert.IsType(resumed.Result).Value; + Assert.False(resumed.IsError); + Assert.Equal(2, result.GetProperty("skipped").GetInt64()); + Assert.Equal(1, result.GetProperty("attempted").GetInt64()); + await fixture.Container.Received(4).PatchItemAsync( + Arg.Any(), + Arg.Any(), + Arg.Any>(), + Arg.Is(options => options.IfMatchEtag != null), + Arg.Any()); + } + + [Fact] + public async Task WhereSelectionHonorsMaxItemsAndBudget() + { + using var fixture = new Fixture(); + fixture.Query("""{"id":"a","__pk0":"t"}""", """{"id":"b","__pk0":"t"}""", """{"id":"c","__pk0":"t"}"""); + + var limited = await fixture.RunAsync("""bulk delete --where "true" --max-items 2 --yes --db db --con items"""); + var result = Assert.IsType(limited.Result).Value; + Assert.True(result.GetProperty("selectionLimited").GetBoolean()); + Assert.Equal(2, result.GetProperty("succeeded").GetInt64()); + + var save = fixture.Directory.PathFor("partial.jsonl"); + var budget = await fixture.RunAsync($"bulk delete --where \"true\" --max-ru 4 --save '{save}' --yes --db db --con items"); + Assert.True(budget.IsError); + Assert.Equal(ShellExitCode.Throttled, budget.ExitCode); + Assert.True(Assert.IsType(budget.Result).Value.GetProperty("budgetExceeded").GetBoolean()); + Assert.False(File.Exists(save)); + Assert.Equal(2, fixture.Writes().Count()); + } + + [Fact] + public async Task SaveNeverOverwritesExistingFiles() + { + using var fixture = new Fixture(); + var save = fixture.Directory.PathFor("existing.jsonl"); + await File.WriteAllTextAsync(save, "keep", TestContext.Current.CancellationToken); + + await Assert.ThrowsAsync(() => fixture.RunAsync($"bulk delete --where \"true\" --dry-run --save '{save}' --db db --con items")); + Assert.Equal("keep", await File.ReadAllTextAsync(save, TestContext.Current.CancellationToken)); + } + + [Fact] + public async Task StatefulJobRetainsOutcomesAndRetriesOnlyFailedOperations() + { + using var fixture = new Fixture(); + await fixture.RunAsync("bulk begin --partition-key a --db db --con items"); + await fixture.RunAsync("""bulk add '{"op":"delete","id":"1"}'"""); + await fixture.RunAsync("""bulk add '[{"op":"patch","id":"2","operations":[{"op":"incr","path":"/n","value":1}]},{"op":"upsert","item":{"id":"3","pk":"a"}}]'"""); + Assert.Equal(3, fixture.Shell.CurrentBulk!.Operations.Count); + + var status = Assert.IsType((await fixture.RunAsync("bulk status")).Result).Value; + Assert.Equal(3, status.GetProperty("operationCount").GetInt32()); + var show = Assert.IsType((await fixture.RunAsync("bulk show")).Result).Value; + Assert.Equal("""["a"]""", show[0].GetProperty("partitionKey").GetRawText()); + + fixture.FailPatchFor("2", HttpStatusCode.PreconditionFailed); + var failed = await fixture.RunAsync("bulk execute --continue-on-error --yes"); + Assert.True(failed.IsError); + Assert.NotNull(fixture.Shell.CurrentBulk); + + fixture.FailPatchFor(null, HttpStatusCode.OK); + var retried = await fixture.RunAsync("bulk commit --yes"); + Assert.False(retried.IsError); + Assert.Equal(2, Assert.IsType(retried.Result).Value.GetProperty("skipped").GetInt64()); + Assert.Null(fixture.Shell.CurrentBulk); + await fixture.Container.Received(1).DeleteItemAsync(Arg.Any(), Arg.Any(), Arg.Any(), Arg.Any()); + } + + [Fact] + public async Task StatefulErrorsMatchBatch() + { + using var fixture = new Fixture(); + await Assert.ThrowsAsync(() => fixture.RunAsync("""bulk add '{"op":"delete","id":"1","partitionKey":"a"}'""")); + await Assert.ThrowsAsync(() => fixture.RunAsync("bulk execute --yes")); + await fixture.RunAsync("bulk begin --db db --con items"); + await Assert.ThrowsAsync(() => fixture.RunAsync("bulk begin --db db --con items")); + await Assert.ThrowsAsync(() => fixture.RunAsync("""bulk add '{"op":"delete","id":"1"}'""")); + Assert.Empty(fixture.Shell.CurrentBulk!.Operations); + await fixture.RunAsync("bulk abort"); + Assert.Null(fixture.Shell.CurrentBulk); + } + + [Fact] + public async Task StatefulExecuteRejectsRecreatedContainer() + { + using var fixture = new Fixture(); + await fixture.RunAsync("bulk begin --db db --con items"); + await fixture.RunAsync("""bulk add '{"op":"delete","id":"1","partitionKey":"a"}'"""); + fixture.Rid = "recreated"; + + await Assert.ThrowsAsync(() => fixture.RunAsync("bulk execute --yes")); + Assert.Empty(fixture.Writes()); + } + + // MCP and help surface. + [Fact] + public void McpSchemaExposesSafeBoundsAndDescribesSubcommands() + { + var factory = new CommandRunner().Commands["bulk"]; + var tool = ToolOperations.GetTool(factory); + var properties = tool.InputSchema.GetProperty("properties"); + + Assert.Equal(16, properties.GetProperty("concurrency").GetProperty("default").GetInt32()); + Assert.Equal(1, properties.GetProperty("concurrency").GetProperty("minimum").GetInt32()); + Assert.Equal(1, properties.GetProperty("max-items").GetProperty("minimum").GetInt32()); + Assert.Equal(0, properties.GetProperty("max-ru").GetProperty("exclusiveMinimum").GetInt32()); + Assert.Contains("subcommand", tool.InputSchema.GetProperty("required").EnumerateArray().Select(value => value.GetString())); + Assert.Contains("patch --where", tool.Description); + Assert.True(factory.McpAnnotation!.Confirmable); + } + + [Fact] + public void HelpExplainsMcpAvailability() + { + var state = HelpCommand.PrintCommandHelp("bulk", new CommandRunner(), plain: true); + var help = Assert.IsType(state.Result).Value; + Assert.Contains("Available through MCP", help.GetProperty("isRestricted").GetString()); + } + + private static JsonElement Json(string value) => JsonSerializer.Deserialize(value); + + private static BulkOperation Operation(long index) => + new(index, "delete", index.ToString(CultureInfo.InvariantCulture), """["a"]""", $$"""{"op":"delete","id":"{{index}}","partitionKey":["a"]}"""); + + private static BulkOutcome Succeeded(long index) => BulkOutcome.Create(Operation(index), BulkOutcome.Succeeded, 1, 200); + + private static Task IgnoreAsync(BulkOutcome outcome) => Task.CompletedTask; + + private static CancellationTokenSource Timeout() + { + var cancellation = CancellationTokenSource.CreateLinkedTokenSource(TestContext.Current.CancellationToken); + cancellation.CancelAfter(TimeSpan.FromSeconds(10)); + return cancellation; + } + + private static async IAsyncEnumerable SourceAsync(int count, [EnumeratorCancellation] CancellationToken token) + { + await Task.CompletedTask; + for (var index = 0; index < count; index++) + { + token.ThrowIfCancellationRequested(); + yield return Operation(index); + } + } + + private sealed class TempDirectory : IDisposable + { + private readonly DirectoryInfo directory = System.IO.Directory.CreateTempSubdirectory("cosmos-bulk-test-"); + + public string PathFor(string name) => Path.Combine(this.directory.FullName, name); + + public void Dispose() => this.directory.Delete(recursive: true); + } + + private sealed class Fixture : IDisposable + { + private string? failingPatchId; + private HttpStatusCode failingStatus = HttpStatusCode.OK; + + public Fixture() + { + var client = Substitute.For(); + var database = Substitute.For(); + client.Endpoint.Returns(new Uri("https://unit-test.documents.azure.com")); + client.GetDatabase("db").Returns(database); + database.GetContainer("items").Returns(this.Container); + this.Shell.State = new ConnectedState(client); + this.Shell.IsInteractiveSession = () => false; + this.Container.ReadContainerStreamAsync(Arg.Any(), Arg.Any()).Returns(_ => + { + var response = new ResponseMessage(HttpStatusCode.OK) + { + Content = new MemoryStream(Encoding.UTF8.GetBytes(JsonSerializer.Serialize(new { _rid = this.Rid, partitionKey = new { paths = SinglePath } }))), + }; + response.Headers.Add("x-ms-request-charge", "2"); + return Task.FromResult(response); + }); + ConfigureWrites(this.Container, HttpStatusCode.OK, 5); + this.Container.PatchItemAsync(default!, default, default!, default, default).ReturnsForAnyArgs(call => + { + if (call.ArgAt(0) == this.failingPatchId) + { + throw new CosmosException("failure", this.failingStatus, 0, "activity", 1); + } + + return Task.FromResult(Response(HttpStatusCode.OK, 5)); + }); + } + + public ShellInterpreter Shell { get; } = ShellInterpreter.CreateInstance(); + + public Container Container { get; } = Substitute.For(); + + public TempDirectory Directory { get; } = new(); + + public string Rid { get; set; } = "rid"; + + public string? LastQuery { get; private set; } + + public static void ConfigureWrites(Container container, HttpStatusCode status, double charge) + { + var response = Response(status, charge); + container.CreateItemAsync(default(JsonElement), default, default, default).ReturnsForAnyArgs(response); + container.UpsertItemAsync(default(JsonElement), default, default, default).ReturnsForAnyArgs(response); + container.ReplaceItemAsync(default(JsonElement), default!, default, default, default).ReturnsForAnyArgs(response); + container.DeleteItemAsync(default!, default, default, default).ReturnsForAnyArgs(response); + container.PatchItemAsync(default!, default, default!, default, default).ReturnsForAnyArgs(response); + } + + public void FailPatchFor(string? id, HttpStatusCode status) + { + this.failingPatchId = id; + this.failingStatus = status; + } + + public IEnumerable Writes() => + this.Container.ReceivedCalls().Where(call => call.GetMethodInfo().Name is "CreateItemAsync" or "UpsertItemAsync" or "ReplaceItemAsync" or "DeleteItemAsync" or "PatchItemAsync"); + + public void Query(params string[] rows) + { + this.Container.GetItemQueryIterator(Arg.Any(), Arg.Any(), Arg.Any()).Returns(call => + { + this.LastQuery = call.ArgAt(0).QueryText; + var index = 0; + var iterator = Substitute.For>(); + iterator.HasMoreResults.Returns(_ => index < rows.Length); + iterator.ReadNextAsync(Arg.Any()).Returns(_ => + { + var row = Json(rows[index++]); + var page = Substitute.For>(); + page.GetEnumerator().Returns(_ => ((IEnumerable)[row]).GetEnumerator()); + page.RequestCharge.Returns(1); + return Task.FromResult(page); + }); + return iterator; + }); + } + + public async Task RunAsync(string commandText) + { + var state = await this.Shell.RunCommandAsync(new CommandState(), commandText, TestContext.Current.CancellationToken); + if (state is ErrorCommandState { Result: null } error) + { + throw error.Exception; + } + + return state; + } + + public void Dispose() + { + this.Shell.Dispose(); + this.Directory.Dispose(); + } + + private static ItemResponse Response(HttpStatusCode status, double charge) + { + var response = Substitute.For>(); + response.StatusCode.Returns(status); + response.RequestCharge.Returns(charge); + return response; + } + } +} diff --git a/CosmosDBShell.Tests/McpConfirmationTests.cs b/CosmosDBShell.Tests/McpConfirmationTests.cs index 238e847b..78d0e81c 100644 --- a/CosmosDBShell.Tests/McpConfirmationTests.cs +++ b/CosmosDBShell.Tests/McpConfirmationTests.cs @@ -16,6 +16,60 @@ namespace CosmosShell.Tests; [Collection(CosmosShell.Tests.Shell.ThemeStateTestCollection.Name)] public class McpConfirmationTests { + [Theory] + [InlineData("run", false, 1, "was not approved")] + [InlineData("run", true, 0, "not connected")] + [InlineData("begin", false, 0, "MCP supports only the stateless 'bulk run'")] + [InlineData("execute", false, 0, "MCP supports only the stateless 'bulk run'")] + public async Task Bulk_McpRequiresConfirmationExceptDryRunAndRejectsStatefulSubcommands(string subcommand, bool dryRun, int expectedPrompts, string expectedText) + { + using var timeout = CancellationTokenSource.CreateLinkedTokenSource(TestContext.Current.CancellationToken); + timeout.CancelAfter(TimeSpan.FromSeconds(10)); + using var host = McpServer.CreateHost(new Program.CosmosShellOptions { McpPort = 0 }); + await host.StartAsync(timeout.Token); + try + { + var prompts = 0; + var address = host.Services.GetRequiredService().Features.Get()!.Addresses.Single(); + var transport = new HttpClientTransport(new HttpClientTransportOptions { Endpoint = new Uri(address.TrimEnd('/') + "/") }); + await using var client = await McpClient.CreateAsync( + transport, + new McpClientOptions + { + Handlers = new McpClientHandlers + { + ElicitationHandler = (_, _) => + { + prompts++; + return ValueTask.FromResult(new ElicitResult { Action = "decline" }); + }, + }, + }, + cancellationToken: timeout.Token); + var arguments = new Dictionary { ["subcommand"] = subcommand, ["yes"] = true }; + if (subcommand == "run") + { + arguments["data"] = "[{\"op\":\"delete\",\"id\":\"1\",\"partitionKey\":\"a\"}]"; + } + + if (dryRun) + { + arguments["dry-run"] = true; + } + + var result = await client.CallToolAsync("bulk", arguments, cancellationToken: timeout.Token); + Assert.Equal(expectedPrompts, prompts); + Assert.True(result.IsError); + var text = Assert.IsType(Assert.Single(result.Content)).Text; + using var document = JsonDocument.Parse(text); + Assert.Contains(expectedText, document.RootElement.GetProperty("error").GetString(), StringComparison.OrdinalIgnoreCase); + } + finally + { + await host.StopAsync(TestContext.Current.CancellationToken); + } + } + [Theory] [InlineData(null, "2026-07-28")] // Stateless request; confirmation uses native multi-round-trip requests. [InlineData("2025-11-25", "2025-11-25")] // Initialize handshake; confirmation is sent over the session. diff --git a/CosmosDBShell.Tests/Shell/CosmosShellPromptTests.cs b/CosmosDBShell.Tests/Shell/CosmosShellPromptTests.cs index 1a08916c..7a47d192 100644 --- a/CosmosDBShell.Tests/Shell/CosmosShellPromptTests.cs +++ b/CosmosDBShell.Tests/Shell/CosmosShellPromptTests.cs @@ -67,4 +67,17 @@ public void GetPromptString_WithActiveBatch_EscapesIndicatorOnce() Assert.Contains(Markup.Escape("[batch:0]"), prompt); Assert.DoesNotContain(Markup.Escape(Markup.Escape("[batch:0]")), prompt); } + + [Fact] + public void GetPromptString_WithActiveBulk_ShowsOperationCount() + { + var shell = ShellInterpreter.CreateInstance(); + shell.State = new DisconnectedState(); + shell.CurrentBulk = new PendingBulkState("TestDatabase", "TestContainer", "rid", ["/pk"], null, null); + shell.CurrentBulk.Operations.Add(new Azure.Data.Cosmos.Shell.Commands.BulkOperation(0, "delete", "1", "[\"a\"]", "{}")); + + var prompt = new CosmosShellPrompt(shell).GetPromptString(); + + Assert.Contains(Markup.Escape("[bulk:1]"), prompt); + } } diff --git a/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BatchOperationParser.cs b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BatchOperationParser.cs index b320ba51..98672f37 100644 --- a/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BatchOperationParser.cs +++ b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BatchOperationParser.cs @@ -52,7 +52,7 @@ public static List Parse(string commandName, string json) } } - private static BatchOperationSpec ParseOne(string commandName, JsonElement element) + internal static BatchOperationSpec ParseOne(string commandName, JsonElement element) { if (element.ValueKind != JsonValueKind.Object) { @@ -157,7 +157,7 @@ private static JsonElement RequireItem(string commandName, JsonElement element, return null; } - private static List ParsePatchOperations(string commandName, JsonElement element) + internal static List ParsePatchOperations(string commandName, JsonElement element) { if (!element.TryGetProperty("operations", out var operationsElement) || operationsElement.ValueKind != JsonValueKind.Array) { diff --git a/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkCommand.cs b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkCommand.cs new file mode 100644 index 00000000..fa607985 --- /dev/null +++ b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkCommand.cs @@ -0,0 +1,949 @@ +// ------------------------------------------------------------ +// Copyright (c) Microsoft Corporation. All rights reserved. +// ------------------------------------------------------------ + +namespace Azure.Data.Cosmos.Shell.Commands; + +using System.Diagnostics; +using System.Globalization; +using System.Net; +using System.Runtime.CompilerServices; +using System.Text; +using System.Text.Encodings.Web; +using System.Text.Json; +using System.Text.Json.Nodes; +using System.Text.Json.Serialization; +using Azure.Data.Cosmos.Shell.Mcp; +using Azure.Data.Cosmos.Shell.Parser; +using Azure.Data.Cosmos.Shell.Util; +using global::Azure; +using global::Azure.Data.Cosmos.Shell.Core; +using global::Azure.Data.Cosmos.Shell.States; +using Spectre.Console; + +[CosmosCommand("bulk")] +[CosmosExample("bulk run '[{\"op\":\"upsert\",\"item\":{\"id\":\"1\",\"pk\":\"a\"}},{\"op\":\"delete\",\"id\":\"2\",\"partitionKey\":\"b\"}]' --concurrency 32 --yes", DescriptionKey = "command-bulk-example-1")] +[CosmosExample("bulk run operations.jsonl --journal operations.journal --yes", DescriptionKey = "command-bulk-example-2")] +[CosmosExample("bulk patch --where 'c.schemaVersion = 1' --operations '[{\"op\":\"set\",\"path\":\"/schemaVersion\",\"value\":2}]' --dry-run --save migration.jsonl", DescriptionKey = "command-bulk-example-3")] +[CosmosExample("bulk delete --where 'c.expired = true' --etag --yes", DescriptionKey = "command-bulk-example-4")] +[CosmosExample("bulk begin --partition-key a", DescriptionKey = "command-bulk-example-5")] +[CosmosExample("bulk add '{\"op\":\"patch\",\"id\":\"3\",\"operations\":[{\"op\":\"set\",\"path\":\"/status\",\"value\":\"done\"}]}'", DescriptionKey = "command-bulk-example-6")] +[CosmosExample("bulk execute --yes", DescriptionKey = "command-bulk-example-7")] +#pragma warning disable SA1118 // Parameter should not span multiple lines +[McpAnnotation( + Title = "Bulk", + Description = @" +Executes many independent write operations with bounded concurrency. Unlike 'batch', operations may target different partition keys, there is no 100-operation limit, and nothing is rolled back: each operation succeeds or fails on its own. + +Subcommands available through MCP: +- 'run ' executes operations. Same operation schema as batch, plus an optional per-operation 'partitionKey' (scalar, or array for hierarchical keys) and 'ifMatch' (ETag). The 'partition-key' argument sets a default for delete and patch operations without one. +- 'patch --where --operations ' patches every item matching a Cosmos SQL predicate over alias c. +- 'delete --where ' deletes every matching item. + +Use 'dry-run' to validate or select without writing; dry runs do not require confirmation. Writes always require user confirmation. Stateful subcommands (begin, add, execute, cancel, status, show) are interactive-shell only. +", + Restricted = true, + Destructive = true, + Confirmable = true)] +#pragma warning restore SA1118 // Parameter should not span multiple lines +internal sealed class BulkCommand : CosmosCommand +{ + internal const int DefaultConcurrency = 16; + + internal static readonly JsonSerializerOptions JsonOptions = new() + { + PropertyNamingPolicy = JsonNamingPolicy.CamelCase, + DefaultIgnoreCondition = JsonIgnoreCondition.WhenWritingNull, + Encoder = JavaScriptEncoder.UnsafeRelaxedJsonEscaping, + }; + + internal static readonly IReadOnlyList McpSubcommands = ["run", "patch", "delete"]; + + private const int SelectionPageSize = 1000; + + private const int StatusPreviewCount = 100; + + private static readonly TimeSpan ProgressInterval = TimeSpan.FromSeconds(2); + + private static readonly Dictionary AllowedOptions = new() + { + ["run"] = ["partition-key", "database", "container", "max-ru", "dry-run", "yes", "continue-on-error", "journal", "retry-uncertain"], + ["patch"] = ["where", "operations", "database", "container", "max-items", "max-ru", "dry-run", "yes", "continue-on-error", "etag", "save", "journal"], + ["delete"] = ["where", "database", "container", "max-items", "max-ru", "dry-run", "yes", "continue-on-error", "etag", "save", "journal"], + ["begin"] = ["partition-key", "database", "container"], + ["add"] = [], + ["execute"] = ["max-ru", "yes", "continue-on-error", "journal", "retry-uncertain"], + ["cancel"] = [], + ["status"] = [], + ["show"] = [], + }; + + private BulkSummary? summary; + + [CosmosParameter("subcommand", RequiredErrorKey = "command-bulk-error-missing_subcommand")] + public string Subcommand { get; init; } = string.Empty; + + [CosmosParameter("data", IsRequired = false)] + public string? Data { get; init; } + + [CosmosOption("where")] + public string? Where { get; init; } + + [CosmosOption("operations")] + public string? Operations { get; init; } + + [CosmosOption("partition-key", "pk")] + public string? PartitionKeyArgument { get; init; } + + [CosmosOption("database", "db")] + public string? Database { get; init; } + + [CosmosOption("container", "con")] + public string? Container { get; init; } + + [CosmosOption("concurrency", "max-parallelism", DefaultValue = DefaultConcurrency)] + public int? Concurrency { get; init; } + + [CosmosOption("max-items")] + public long? MaxItems { get; init; } + + [CosmosOption("max-ru")] + public double? MaxRu { get; init; } + + [CosmosOption("dry-run")] + public bool? DryRun { get; init; } + + [CosmosOption("yes", "y")] + public bool? Yes { get; init; } + + [CosmosOption("continue-on-error")] + public bool? ContinueOnError { get; init; } + + [CosmosOption("etag")] + public bool? ETag { get; init; } + + [CosmosOption("save")] + public string? Save { get; init; } + + [CosmosOption("journal")] + public string? Journal { get; init; } + + [CosmosOption("retry-uncertain")] + public bool? RetryUncertain { get; init; } + + internal bool ConfirmationApproved { get; set; } + + internal static string NormalizeSubcommand(string? subcommand) + { + return subcommand?.Trim().ToLowerInvariant() switch + { + "exec" or "commit" => "execute", + "abort" => "cancel", + var value => value ?? string.Empty, + }; + } + + public override async Task ExecuteAsync(ShellInterpreter shell, CommandState commandState, string commandText, CancellationToken token) + { + var subcommand = NormalizeSubcommand(this.Subcommand); + this.Validate(subcommand); + try + { + return subcommand switch + { + "run" => await this.RunAsync(shell, commandState, token), + "patch" or "delete" => await this.RunSelectionAsync(shell, subcommand, token), + "begin" => await this.BeginAsync(shell, token), + "add" => await this.AddAsync(shell, commandState, token), + "execute" => await this.ExecutePendingAsync(shell, token), + "cancel" => Cancel(shell), + "status" => Status(shell), + _ => Show(shell), + }; + } + catch (Exception) when (this.summary is not null) + { + RequestChargeContext.Record(this.summary.RequestCharge); + throw; + } + } + + internal static async Task WriteAsync(Container container, BulkOperation operation, CancellationToken token) + { + try + { + using var document = JsonDocument.Parse(operation.Json); + var root = document.RootElement; + var key = BulkOperationNormalizer.ParseKey(operation.PartitionKey); + var ifMatch = root.TryGetProperty("ifMatch", out var tag) ? tag.GetString() : null; + var options = new ItemRequestOptions { IfMatchEtag = ifMatch, EnableContentResponseOnWrite = false }; + var item = root.TryGetProperty("item", out var value) ? value.Clone() : default; + ItemResponse response = operation.Op switch + { + "create" => await container.CreateItemAsync(item, key, options, token), + "upsert" => await container.UpsertItemAsync(item, key, options, token), + "replace" => await container.ReplaceItemAsync(item, operation.Id, key, options, token), + "delete" => await container.DeleteItemAsync(operation.Id, key, options, token), + "patch" => await container.PatchItemAsync( + operation.Id, + key, + BatchOperationParser.ParsePatchOperations("bulk", root), + new PatchItemRequestOptions { IfMatchEtag = ifMatch, EnableContentResponseOnWrite = false }, + token), + _ => throw new InvalidOperationException($"Unsupported bulk operation '{operation.Op}'."), + }; + + var status = (int)response.StatusCode; + return BulkOutcome.Create( + operation, + status is >= 200 and < 300 ? BulkOutcome.Succeeded : BulkOutcome.Failed, + response.RequestCharge, + status); + } + catch (CosmosException ex) + { + // Timeouts and server errors can occur after the service applied the write. + var uncertain = ex.StatusCode == HttpStatusCode.RequestTimeout || (int)ex.StatusCode >= 500; + return BulkOutcome.Create( + operation, + uncertain ? BulkOutcome.Uncertain : BulkOutcome.Failed, + ex.RequestCharge, + (int)ex.StatusCode, + CommandException.GetDisplayMessage(ex)); + } + } + + internal static string BuildSelectionQuery(string where, IReadOnlyList partitionKeyPaths) + { + var query = new StringBuilder("SELECT c.id, c._etag"); + for (var i = 0; i < partitionKeyPaths.Count; i++) + { + query.Append(", c"); + foreach (var segment in partitionKeyPaths[i].Split('/', StringSplitOptions.RemoveEmptyEntries)) + { + query.Append('[').Append(JsonSerializer.Serialize(segment, JsonOptions)).Append(']'); + } + + query.Append(" AS __pk").Append(i.ToString(CultureInfo.InvariantCulture)); + } + + return query.Append(" FROM c WHERE (").Append(where).Append(')').ToString(); + } + + internal void Validate(string subcommand) + { + if (subcommand.Length == 0) + { + throw Error("command-bulk-error-missing_subcommand"); + } + + if (!AllowedOptions.TryGetValue(subcommand, out var allowed)) + { + throw new CommandException("bulk", MessageService.GetArgsString("command-bulk-error-invalid_subcommand", "subcommand", subcommand)); + } + + foreach (var (name, supplied) in this.GetSuppliedOptions()) + { + if (supplied && !allowed.Contains(name)) + { + throw new CommandException("bulk", MessageService.GetArgsString("command-bulk-error-option_not_supported", "option", name, "subcommand", subcommand)); + } + } + + if (this.Data != null && subcommand is not ("run" or "add")) + { + throw new CommandException("bulk", MessageService.GetArgsString("command-bulk-error-data_not_supported", "subcommand", subcommand)); + } + + if (this.Concurrency is < 1 || this.MaxItems is < 1 || (this.MaxRu is { } ru && (!double.IsFinite(ru) || ru <= 0))) + { + throw Error("command-bulk-error-invalid_number"); + } + + if (subcommand is "patch" or "delete" && string.IsNullOrWhiteSpace(this.Where)) + { + throw new CommandException("bulk", MessageService.GetArgsString("command-bulk-error-missing_where", "subcommand", subcommand)); + } + + if (subcommand == "patch") + { + if (string.IsNullOrWhiteSpace(this.Operations)) + { + throw Error("command-bulk-error-missing_operations"); + } + + try + { + BulkOperationNormalizer.ValidatePatchOperations(JsonSerializer.Deserialize(this.Operations)); + } + catch (JsonException ex) + { + throw new CommandException("bulk", MessageService.GetArgsString("command-batch-error-invalid_json", "message", ex.Message), ex); + } + } + + if (this.DryRun == true && (this.Journal != null || this.RetryUncertain == true)) + { + throw new CommandException("bulk", MessageService.GetArgsString("command-bulk-error-dry_run_option", "option", this.Journal != null ? "journal" : "retry-uncertain")); + } + + if (this.Journal != null && subcommand is "patch" or "delete" && this.Save == null) + { + throw Error("command-bulk-error-journal_requires_save"); + } + + if (this.RetryUncertain == true && subcommand == "run" && this.Journal == null) + { + throw Error("command-bulk-error-retry_requires_journal"); + } + } + + private static CommandException Error(string key) => new("bulk", MessageService.GetString(key)); + + private static ConnectedState RequireConnected(ShellInterpreter shell) + { + return shell.State as ConnectedState ?? throw new NotConnectedException("bulk"); + } + + private static bool CanPrompt(ShellInterpreter shell) + { + return !shell.IsMachineMode + && shell.IsInteractiveSession() + && string.IsNullOrEmpty(shell.CurrentScriptFileName) + && shell.Options?.ExecuteAndQuit == null + && shell.Options?.ExecuteAndContinue == null; + } + + private static async IAsyncEnumerable ReadOperationsAsync(string data, [EnumeratorCancellation] CancellationToken token) + { + var trimmed = data.TrimStart(); + if (trimmed.StartsWith('[') || trimmed.StartsWith('{')) + { + using var document = ParseInline(data); + if (document.RootElement.ValueKind == JsonValueKind.Array) + { + foreach (var element in document.RootElement.EnumerateArray()) + { + yield return element.Clone(); + } + } + else + { + yield return document.RootElement.Clone(); + } + + yield break; + } + + if (!File.Exists(data)) + { + throw new CommandException("bulk", MessageService.GetArgsString("command-bulk-error-file_not_found", "file", data)); + } + + await using var stream = new FileStream(data, FileMode.Open, FileAccess.Read, FileShare.Read, 65536, useAsync: true); + var isArray = await StartsWithArrayAsync(stream, token); + stream.Position = 0; + if (isArray) + { + await foreach (var (_, element) in ImportCommand.EnumerateArrayAsync(stream, token)) + { + yield return element; + } + } + else + { + using var reader = new StreamReader(stream, Encoding.UTF8, detectEncodingFromByteOrderMarks: true); + await foreach (var (_, element) in ImportCommand.EnumerateJsonLinesAsync(reader, token)) + { + yield return element; + } + } + } + + private static JsonDocument ParseInline(string data) + { + try + { + return JsonDocument.Parse(data); + } + catch (JsonException ex) + { + throw new CommandException("bulk", MessageService.GetArgsString("command-batch-error-invalid_json", "message", ex.Message), ex); + } + } + + private static async Task StartsWithArrayAsync(Stream stream, CancellationToken token) + { + var buffer = new byte[4096]; + int read; + while ((read = await stream.ReadAsync(buffer, token)) > 0) + { + for (var i = 0; i < read; i++) + { + var value = buffer[i]; + if (value is (byte)' ' or (byte)'\t' or (byte)'\r' or (byte)'\n' or 0xEF or 0xBB or 0xBF) + { + continue; + } + + return value == (byte)'['; + } + } + + return false; + } + + private static async IAsyncEnumerable EnumerateAsync(IReadOnlyList operations) + { + foreach (var operation in operations) + { + yield return operation; + } + + await Task.CompletedTask; + } + + private static StreamWriter OpenNewFile(string path) + { + var options = new FileStreamOptions { Mode = FileMode.CreateNew, Access = FileAccess.Write, Share = FileShare.None }; + if (!OperatingSystem.IsWindows()) + { + options.UnixCreateMode = UnixFileMode.UserRead | UnixFileMode.UserWrite; + } + + try + { + return new StreamWriter(new FileStream(path, options), new UTF8Encoding(false)) { NewLine = "\n" }; + } + catch (IOException ex) when (File.Exists(path)) + { + throw new CommandException("bulk", MessageService.GetArgsString("command-bulk-error-save_exists", "file", path), ex); + } + } + + private static JsonElement CreateSelectedOperation(string subcommand, JsonElement row, int componentCount, JsonNode? patchOperations, bool useETag) + { + if (!row.TryGetProperty("id", out var id) || id.ValueKind != JsonValueKind.String) + { + throw Error("command-bulk-error-selected_id"); + } + + var key = new JsonArray(); + for (var i = 0; i < componentCount; i++) + { + key.Add(row.TryGetProperty($"__pk{i}", out var component) ? JsonNode.Parse(component.GetRawText()) : new JsonObject()); + } + + var operation = new JsonObject + { + ["op"] = subcommand, + ["id"] = id.GetString(), + ["partitionKey"] = key, + }; + + if (useETag) + { + if (!row.TryGetProperty("_etag", out var etag) || etag.ValueKind != JsonValueKind.String) + { + throw Error("command-bulk-error-selected_etag"); + } + + operation["ifMatch"] = etag.GetString(); + } + + if (patchOperations != null) + { + operation["operations"] = patchOperations.DeepClone(); + } + + return JsonSerializer.SerializeToElement(operation); + } + + private static CommandState CreateResult(BulkSummary summary) + { + var result = new ShellJson(JsonSerializer.SerializeToElement(summary, JsonOptions)); + var charge = summary.RequestCharge.ToString("F2", CultureInfo.InvariantCulture); + if (summary.Success) + { + var message = summary.DryRun + ? MessageService.GetArgsString("command-bulk-dry-run", "count", summary.OperationCount, "charge", charge) + : MessageService.GetArgsString("command-bulk-success", "count", summary.OperationCount, "succeeded", summary.Succeeded, "skipped", summary.Skipped, "charge", charge); + return new CommandState + { + Result = result, + RequestCharge = summary.RequestCharge, + RenderUser = () => ShellInterpreter.WriteLine(message), + }; + } + + var error = summary.BudgetExceeded + ? MessageService.GetArgsString("command-bulk-error-budget", "succeeded", summary.Succeeded, "count", summary.OperationCount, "charge", charge) + : MessageService.GetArgsString("command-bulk-error-failed", "succeeded", summary.Succeeded, "failed", summary.Failed, "uncertain", summary.Uncertain, "charge", charge); + var state = new StructuredErrorCommandState( + new CommandException("bulk", error, new RequestFailedException(summary.BudgetExceeded ? 429 : 400, error)), + result) + { + RequestCharge = summary.RequestCharge, + }; + state.RenderUser = () => + { + ShellInterpreter.WriteLine(error); + foreach (var failure in summary.Errors.Take(5)) + { + ShellInterpreter.WriteLine(MessageService.GetArgsString( + "command-bulk-error-item", + "index", + failure.Index + 1, + "op", + failure.Op, + "id", + failure.Id, + "status", + failure.StatusCode?.ToString(CultureInfo.InvariantCulture) ?? failure.Status, + "message", + failure.Error ?? failure.Status)); + } + }; + return state; + } + + private static CommandState Cancel(ShellInterpreter shell) + { + var bulk = shell.CurrentBulk ?? throw Error("command-bulk-error-not_active"); + shell.CurrentBulk = null; + var message = MessageService.GetArgsString("command-bulk-cancelled", "count", bulk.Operations.Count); + return new CommandState { RenderUser = () => ShellInterpreter.WriteLine(message) }; + } + + private static CommandState Status(ShellInterpreter shell) + { + var bulk = shell.CurrentBulk; + JsonObject root; + Action renderUser; + if (bulk is null) + { + root = new JsonObject { ["active"] = false }; + renderUser = () => ShellInterpreter.WriteLine(MessageService.GetString("command-bulk-status-inactive")); + } + else + { + var operations = new JsonArray(); + foreach (var operation in bulk.Operations.Take(StatusPreviewCount)) + { + operations.Add(new JsonObject { ["op"] = operation.Op, ["id"] = operation.Id }); + } + + root = new JsonObject + { + ["active"] = true, + ["database"] = bulk.DatabaseName, + ["container"] = bulk.ContainerName, + ["partitionKey"] = bulk.PartitionKeyArgument, + ["operationCount"] = bulk.Operations.Count, + ["succeeded"] = bulk.Outcomes.Values.Count(outcome => outcome.Status == BulkOutcome.Succeeded), + ["failed"] = bulk.Outcomes.Values.Count(outcome => outcome.Status == BulkOutcome.Failed), + ["uncertain"] = bulk.Outcomes.Values.Count(outcome => outcome.Status is BulkOutcome.Started or BulkOutcome.Uncertain), + ["operations"] = operations, + }; + renderUser = () => RenderStatus(bulk); + } + + using var document = JsonDocument.Parse(root.ToJsonString()); + return new CommandState { Result = new ShellJson(document.RootElement.Clone()), RenderUser = renderUser }; + } + + private static void RenderStatus(PendingBulkState bulk) + { + var details = new Table().HideHeaders(); + details.AddColumn(string.Empty); + details.AddColumn(string.Empty); + void AddDetail(string labelKey, string value) => + details.AddRow( + Theme.FormatHelpName(Markup.Escape(MessageService.GetString(labelKey))), + Theme.FormatTableValue(Markup.Escape(value))); + + AddDetail("command-batch-status-target", $"{bulk.DatabaseName}/{bulk.ContainerName}"); + if (bulk.PartitionKeyArgument != null) + { + AddDetail("command-batch-status-partition-key", bulk.PartitionKeyArgument); + } + + AddDetail("command-batch-status-operation-count", bulk.Operations.Count.ToString(CultureInfo.InvariantCulture)); + AnsiConsole.Write(details); + if (bulk.Operations.Count == 0) + { + return; + } + + var operations = new Table(); + operations.AddColumn(Theme.FormatSectionHeader(MessageService.GetString("command-batch-status-column-index"))); + operations.AddColumn(Theme.FormatSectionHeader(MessageService.GetString("command-batch-status-column-operation"))); + operations.AddColumn(Theme.FormatSectionHeader(MessageService.GetString("command-batch-status-column-id"))); + foreach (var operation in bulk.Operations.Take(StatusPreviewCount)) + { + operations.AddRow( + Theme.FormatTableValue((operation.Index + 1).ToString(CultureInfo.InvariantCulture)), + Theme.FormatTableValue(Markup.Escape(operation.Op)), + Theme.FormatTableValue(Markup.Escape(operation.Id))); + } + + AnsiConsole.Write(operations); + if (bulk.Operations.Count > StatusPreviewCount) + { + ShellInterpreter.WriteLine(MessageService.GetArgsString("command-bulk-status-more", "count", bulk.Operations.Count - StatusPreviewCount)); + } + } + + private static CommandState Show(ShellInterpreter shell) + { + var operations = new JsonArray(); + if (shell.CurrentBulk is { } bulk) + { + foreach (var operation in bulk.Operations) + { + operations.Add(JsonNode.Parse(operation.Json)); + } + } + + using var document = JsonDocument.Parse(operations.ToJsonString(JsonOptions)); + return new CommandState { Result = new ShellJson(document.RootElement.Clone()) }; + } + + private IEnumerable<(string Name, bool Supplied)> GetSuppliedOptions() + { + return + [ + ("where", this.Where != null), + ("operations", this.Operations != null), + ("partition-key", this.PartitionKeyArgument != null), + ("database", this.Database != null), + ("container", this.Container != null), + ("max-items", this.MaxItems != null), + ("max-ru", this.MaxRu != null), + ("dry-run", this.DryRun == true), + ("yes", this.Yes == true), + ("continue-on-error", this.ContinueOnError == true), + ("etag", this.ETag == true), + ("save", this.Save != null), + ("journal", this.Journal != null), + ("retry-uncertain", this.RetryUncertain == true), + ]; + } + + private string ResolveData(CommandState commandState) + { + var data = this.Data ?? commandState.Result?.ConvertShellObject(DataType.Text) as string; + return string.IsNullOrWhiteSpace(data) ? throw Error("command-bulk-error-missing_data") : data; + } + + private void RequireApprovalAvailable(ShellInterpreter shell) + { + if (this.DryRun != true && !this.ConfirmationApproved && this.Yes != true && !CanPrompt(shell)) + { + throw Error("command-bulk-error-confirm_required"); + } + } + + private void Confirm(ShellInterpreter shell, long count, Target target) + { + if (this.ConfirmationApproved || this.Yes == true) + { + return; + } + + if (!CanPrompt(shell)) + { + throw Error("command-bulk-error-confirm_required"); + } + + ShellInterpreter.WriteLine(MessageService.GetArgsString( + "command-bulk-confirm-summary", + "count", + count, + "target", + $"{target.DatabaseName}/{target.ContainerName}")); + if (!ShellInterpreter.Confirm("command-bulk-confirm")) + { + throw Error("command-bulk-error-declined"); + } + } + + private async Task ResolveTargetAsync(ConnectedState connected, State state, string? database, string? container, CancellationToken token) + { + var (databaseName, containerName, reference) = ResolveContainerReference(connected.Client, state, database, container, "bulk"); + using var response = await reference.ReadContainerStreamAsync(cancellationToken: token); + this.summary = new BulkSummary { DryRun = this.DryRun == true, RequestCharge = response.Headers.RequestCharge }; + if (!response.IsSuccessStatusCode) + { + var message = MessageService.GetArgsString( + "command-bulk-error-container", + "target", + $"{databaseName}/{containerName}", + "status", + (int)response.StatusCode); + throw new CommandException("bulk", message, new RequestFailedException((int)response.StatusCode, message)); + } + + using var metadata = await JsonDocument.ParseAsync(response.Content, cancellationToken: token); + var root = metadata.RootElement; + string[] paths = root.TryGetProperty("partitionKey", out var definition) && definition.TryGetProperty("paths", out var values) + ? values.EnumerateArray().Select(path => path.GetString()!).ToArray() + : []; + return new Target(databaseName, containerName, reference, connected.Client.Endpoint.ToString(), root.GetProperty("_rid").GetString()!, paths); + } + + private async Task RunAsync(ShellInterpreter shell, CommandState commandState, CancellationToken token) + { + var data = this.ResolveData(commandState); + var connected = RequireConnected(shell); + this.RequireApprovalAvailable(shell); + var target = await this.ResolveTargetAsync(connected, shell.State, this.Database, this.Container, token); + var defaultKey = this.PartitionKeyArgument is null + ? null + : BulkOperationNormalizer.CanonicalKeyFromArgument(this.PartitionKeyArgument, target.PartitionKeyPaths.Count); + + using var spool = BulkSpool.Create(); + await foreach (var element in ReadOperationsAsync(data, token)) + { + var operation = BulkOperationNormalizer.Normalize(element, target.PartitionKeyPaths, defaultKey, spool.Count); + await spool.AddAsync(operation.Json); + } + + if (spool.Count == 0) + { + throw Error("command-bulk-error-empty"); + } + + return await this.ExecuteOperationsAsync(shell, target, spool.Count, spool.GetHash(), spool.ReadAsync(token), null, token); + } + + private async Task RunSelectionAsync(ShellInterpreter shell, string subcommand, CancellationToken token) + { + var connected = RequireConnected(shell); + this.RequireApprovalAvailable(shell); + var target = await this.ResolveTargetAsync(connected, shell.State, this.Database, this.Container, token); + var summary = this.summary!; + var patchOperations = subcommand == "patch" ? JsonNode.Parse(this.Operations!) : null; + + using var spool = BulkSpool.Create(); + var savePath = this.Save is null ? null : Path.GetFullPath(this.Save); + var save = savePath is null ? null : OpenNewFile(savePath); + var complete = false; + try + { + var query = BuildSelectionQuery(this.Where!, target.PartitionKeyPaths); + using var iterator = target.Container.GetItemQueryIterator( + new QueryDefinition(query), + continuationToken: null, + new QueryRequestOptions { MaxItemCount = SelectionPageSize }); + while (iterator.HasMoreResults && !summary.SelectionLimited) + { + if (this.MaxRu is { } budget && summary.RequestCharge >= budget) + { + summary.ResultIncomplete = true; + summary.BudgetExceeded = true; + break; + } + + var page = await iterator.ReadNextAsync(token); + summary.RequestCharge += page.RequestCharge; + foreach (var row in page) + { + if (this.MaxItems is { } max && spool.Count >= max) + { + summary.SelectionLimited = true; + break; + } + + var element = CreateSelectedOperation(subcommand, row, target.PartitionKeyPaths.Count, patchOperations, this.ETag == true); + var operation = BulkOperationNormalizer.Normalize(element, target.PartitionKeyPaths, null, spool.Count); + await spool.AddAsync(operation.Json); + if (save != null) + { + await save.WriteLineAsync(operation.Json); + } + } + + if (this.MaxItems is { } limit && spool.Count >= limit && iterator.HasMoreResults) + { + summary.SelectionLimited = true; + } + } + + complete = !summary.ResultIncomplete; + } + finally + { + if (save != null) + { + await save.DisposeAsync(); + if (!complete) + { + File.Delete(savePath!); + } + } + } + + if (!complete) + { + summary.OperationCount = spool.Count; + return CreateResult(summary); + } + + return await this.ExecuteOperationsAsync(shell, target, spool.Count, spool.GetHash(), spool.ReadAsync(token), null, token); + } + + private async Task BeginAsync(ShellInterpreter shell, CancellationToken token) + { + if (shell.CurrentBulk is not null) + { + throw Error("command-bulk-error-already_active"); + } + + var connected = RequireConnected(shell); + var target = await this.ResolveTargetAsync(connected, shell.State, this.Database, this.Container, token); + var defaultKey = this.PartitionKeyArgument is null + ? null + : BulkOperationNormalizer.CanonicalKeyFromArgument(this.PartitionKeyArgument, target.PartitionKeyPaths.Count); + shell.CurrentBulk = new PendingBulkState( + target.DatabaseName, + target.ContainerName, + target.ContainerRid, + target.PartitionKeyPaths, + this.PartitionKeyArgument, + defaultKey); + + var message = MessageService.GetArgsString("command-bulk-begun", "database", target.DatabaseName, "container", target.ContainerName); + return new CommandState + { + RequestCharge = this.summary!.RequestCharge, + RenderUser = () => ShellInterpreter.WriteLine(message), + }; + } + + private async Task AddAsync(ShellInterpreter shell, CommandState commandState, CancellationToken token) + { + var bulk = shell.CurrentBulk ?? throw Error("command-bulk-error-not_active"); + var data = this.ResolveData(commandState); + var added = new List(); + await foreach (var element in ReadOperationsAsync(data, token)) + { + added.Add(BulkOperationNormalizer.Normalize(element, bulk.PartitionKeyPaths, bulk.DefaultPartitionKey, bulk.Operations.Count + added.Count)); + } + + if (added.Count == 0) + { + throw Error("command-bulk-error-empty"); + } + + bulk.Operations.AddRange(added); + var message = MessageService.GetArgsString("command-bulk-added", "count", added.Count, "total", bulk.Operations.Count); + return new CommandState { RenderUser = () => ShellInterpreter.WriteLine(message) }; + } + + private async Task ExecutePendingAsync(ShellInterpreter shell, CancellationToken token) + { + var bulk = shell.CurrentBulk ?? throw Error("command-bulk-error-not_active"); + if (bulk.Operations.Count == 0) + { + throw Error("command-bulk-error-empty"); + } + + var connected = RequireConnected(shell); + this.RequireApprovalAvailable(shell); + var target = await this.ResolveTargetAsync(connected, shell.State, bulk.DatabaseName, bulk.ContainerName, token); + if (target.ContainerRid != bulk.ContainerRid) + { + throw Error("command-bulk-error-container_changed"); + } + + var hash = BulkSpool.ComputeHash(bulk.Operations.Select(operation => operation.Json)); + var result = await this.ExecuteOperationsAsync(shell, target, bulk.Operations.Count, hash, EnumerateAsync(bulk.Operations), bulk.Outcomes, token); + if (!result.IsError) + { + shell.CurrentBulk = null; + } + + return result; + } + + private async Task ExecuteOperationsAsync( + ShellInterpreter shell, + Target target, + long count, + string hash, + IAsyncEnumerable operations, + Dictionary? memory, + CancellationToken token) + { + var summary = this.summary!; + summary.OperationCount = count; + if (summary.DryRun || count == 0) + { + return CreateResult(summary); + } + + this.Confirm(shell, count, target); + using var journal = this.Journal is null + ? null + : BulkJournal.Open(this.Journal, target.Endpoint, target.ContainerRid, hash, count); + var previous = new Dictionary(); + foreach (var source in new[] { journal?.Previous, memory }) + { + foreach (var (index, outcome) in source ?? Enumerable.Empty>()) + { + previous[index] = outcome; + } + } + + var progress = Stopwatch.StartNew(); + async Task RecordAsync(BulkOutcome outcome) + { + if (memory != null) + { + memory[outcome.Index] = outcome; + } + + if (journal != null) + { + await journal.RecordAsync(outcome); + } + + if (outcome.Status != BulkOutcome.Started && progress.Elapsed >= ProgressInterval && !shell.IsMachineMode) + { + progress.Restart(); + ShellInterpreter.WriteLine(MessageService.GetArgsString( + "command-bulk-progress", + "processed", + summary.Processed, + "total", + count, + "failed", + summary.Failed, + "charge", + summary.RequestCharge.ToString("F2", CultureInfo.InvariantCulture))); + } + } + + await BulkExecutor.ExecuteAsync( + operations, + (operation, cancellation) => WriteAsync(target.Container, operation, cancellation), + RecordAsync, + previous, + summary, + this.Concurrency ?? DefaultConcurrency, + this.MaxRu, + this.ContinueOnError == true, + this.RetryUncertain == true, + token); + return CreateResult(summary); + } + + private sealed record Target( + string DatabaseName, + string ContainerName, + Container Container, + string Endpoint, + string ContainerRid, + IReadOnlyList PartitionKeyPaths); +} diff --git a/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkExecutor.cs b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkExecutor.cs new file mode 100644 index 00000000..4a783cef --- /dev/null +++ b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkExecutor.cs @@ -0,0 +1,141 @@ +// ------------------------------------------------------------ +// Copyright (c) Microsoft Corporation. All rights reserved. +// ------------------------------------------------------------ + +namespace Azure.Data.Cosmos.Shell.Commands; + +using System.Runtime.ExceptionServices; + +internal static class BulkExecutor +{ + public static async Task ExecuteAsync( + IAsyncEnumerable source, + Func> write, + Func record, + IReadOnlyDictionary previous, + BulkSummary summary, + int concurrency, + double? maxRu, + bool continueOnError, + bool retryUncertain, + CancellationToken token) + { + ArgumentOutOfRangeException.ThrowIfLessThan(concurrency, 1); + var pending = new List<(BulkOperation Operation, PartitionKey Key, Task Task)>(); + var newFailures = 0; + ExceptionDispatchInfo? failure = null; + + async Task SettleAsync(int index) + { + var (operation, _, task) = pending[index]; + pending.RemoveAt(index); + BulkOutcome result; + try + { + result = await task; + } + catch (Exception ex) + { + // The request may or may not have reached the service. + failure ??= ExceptionDispatchInfo.Capture(ex); + result = BulkOutcome.Create(operation, BulkOutcome.Uncertain, error: ex.Message); + } + + summary.RequestCharge += result.RequestCharge; + summary.Processed++; + if (result.Status == BulkOutcome.Succeeded) + { + summary.Succeeded++; + } + else + { + newFailures++; + summary.AddFailure(result); + } + + await record(result); + } + + async Task SettleCompletedAsync() + { + for (var i = pending.Count - 1; i >= 0; i--) + { + if (pending[i].Task.IsCompleted) + { + await SettleAsync(i); + } + } + } + + try + { + await foreach (var operation in source.WithCancellation(token)) + { + var key = BulkOperationNormalizer.ParseKey(operation.PartitionKey); + await SettleCompletedAsync(); + + // Operations on the same item keep their list order. + while (pending.Count >= concurrency || pending.Any(entry => entry.Operation.Id == operation.Id && entry.Key.Equals(key))) + { +#pragma warning disable VSTHRD003 // All pending tasks were started by this invocation. + await Task.WhenAny(pending.Select(entry => entry.Task)).WaitAsync(token); +#pragma warning restore VSTHRD003 + await SettleCompletedAsync(); + } + + var budgetExhausted = maxRu.HasValue && summary.RequestCharge >= maxRu.Value; + if (failure != null || (newFailures > 0 && !continueOnError) || budgetExhausted) + { + summary.ResultIncomplete = true; + summary.BudgetExceeded = budgetExhausted; + break; + } + + token.ThrowIfCancellationRequested(); + if (previous.TryGetValue(operation.Index, out var old)) + { + if (old.Status == BulkOutcome.Succeeded) + { + summary.Skipped++; + summary.Processed++; + continue; + } + + if (old.Status is BulkOutcome.Started or BulkOutcome.Uncertain && !retryUncertain) + { + summary.Skipped++; + summary.Processed++; + summary.AddFailure(old with { Status = BulkOutcome.Uncertain }); + continue; + } + } + + await record(BulkOutcome.Create(operation, BulkOutcome.Started)); + summary.Attempted++; + pending.Add((operation, key, write(operation, token))); + } + } + catch (Exception ex) + { + failure ??= ExceptionDispatchInfo.Capture(ex); + } + finally + { + // Observe every started write even when reading, cancellation, or journal I/O fails. + while (pending.Count > 0) + { + try + { + await SettleAsync(0); + } + catch (Exception ex) + { + failure ??= ExceptionDispatchInfo.Capture(ex); + } + } + } + + failure?.Throw(); + token.ThrowIfCancellationRequested(); + } +} diff --git a/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkJournal.cs b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkJournal.cs new file mode 100644 index 00000000..a1d18940 --- /dev/null +++ b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkJournal.cs @@ -0,0 +1,165 @@ +// ------------------------------------------------------------ +// Copyright (c) Microsoft Corporation. All rights reserved. +// ------------------------------------------------------------ + +namespace Azure.Data.Cosmos.Shell.Commands; + +using System.Text; +using System.Text.Json; +using Azure.Data.Cosmos.Shell.Core; +using Azure.Data.Cosmos.Shell.Util; + +/// +/// A durable, append-only JSON Lines log of write intents and outcomes for one operation list. +/// The intent is recorded before each write is sent, so a resumed run never mistakes an +/// interrupted write for one that was not attempted. +/// +internal sealed class BulkJournal : IDisposable +{ + private const int Version = 1; + + private readonly FileStream stream; + + private BulkJournal(FileStream stream, Dictionary previous) + { + this.stream = stream; + this.Previous = previous; + } + + public IReadOnlyDictionary Previous { get; } + + public static BulkJournal Open(string path, string endpoint, string containerRid, string hash, long operationCount) + { + var fullPath = Path.GetFullPath(path); + var options = new FileStreamOptions + { + Mode = FileMode.OpenOrCreate, + Access = FileAccess.ReadWrite, + Share = FileShare.Read, + Options = FileOptions.WriteThrough, + }; + if (!OperatingSystem.IsWindows()) + { + options.UnixCreateMode = UnixFileMode.UserRead | UnixFileMode.UserWrite; + } + + FileStream stream; + try + { + stream = new FileStream(fullPath, options); + } + catch (IOException ex) + { + throw new CommandException("bulk", MessageService.GetArgsString("command-bulk-error-journal_open", "file", fullPath, "message", ex.Message), ex); + } + + try + { + var header = new Header(Version, endpoint, containerRid, hash, operationCount); + TrimIncompleteLastLine(stream); + var previous = new Dictionary(); + if (stream.Length == 0) + { + Append(stream, JsonSerializer.Serialize(header, BulkCommand.JsonOptions)); + } + else + { + ReadExisting(stream, header, previous); + stream.Seek(0, SeekOrigin.End); + } + + return new BulkJournal(stream, previous); + } + catch + { + stream.Dispose(); + throw; + } + } + + public async Task RecordAsync(BulkOutcome outcome) + { + var bytes = Encoding.UTF8.GetBytes(JsonSerializer.Serialize(outcome, BulkCommand.JsonOptions) + "\n"); + await this.stream.WriteAsync(bytes); + await this.stream.FlushAsync(); + } + + public void Dispose() => this.stream.Dispose(); + + private static void Append(FileStream stream, string line) + { + var bytes = Encoding.UTF8.GetBytes(line + "\n"); + stream.Write(bytes); + stream.Flush(); + } + + // A crash can leave a partially written final line. Its write had not been sent (intent) or + // its outcome is unknown (the earlier intent remains), so dropping it never hides an attempt. + private static void TrimIncompleteLastLine(FileStream stream) + { + var length = stream.Length; + if (length == 0) + { + return; + } + + var buffer = new byte[4096]; + var end = length; + while (end > 0) + { + var start = Math.Max(0, end - buffer.Length); + stream.Position = start; + stream.ReadExactly(buffer, 0, (int)(end - start)); + for (var i = (int)(end - start) - 1; i >= 0; i--) + { + if (buffer[i] == (byte)'\n') + { + var keep = start + i + 1; + if (keep != length) + { + stream.SetLength(keep); + } + + return; + } + } + + end = start; + } + + stream.SetLength(0); + } + + private static void ReadExisting(FileStream stream, Header expected, Dictionary previous) + { + stream.Position = 0; + using var reader = new StreamReader(stream, new UTF8Encoding(false), false, 65536, leaveOpen: true); + try + { + var header = JsonSerializer.Deserialize
(reader.ReadLine() ?? string.Empty, BulkCommand.JsonOptions); + if (header != expected) + { + throw new CommandException("bulk", MessageService.GetString("command-bulk-error-journal_mismatch")); + } + + while (reader.ReadLine() is { } line) + { + var outcome = JsonSerializer.Deserialize(line, BulkCommand.JsonOptions); + if (outcome is null || outcome.Index < 0 || outcome.Index >= expected.OperationCount || !BulkOutcome.IsKnownStatus(outcome.Status)) + { + throw Invalid(null); + } + + previous[outcome.Index] = outcome; + } + } + catch (JsonException ex) + { + throw Invalid(ex); + } + } + + private static CommandException Invalid(Exception? inner) => new("bulk", MessageService.GetString("command-bulk-error-journal_invalid"), inner); + + private sealed record Header(int Version, string Endpoint, string ContainerRid, string Hash, long OperationCount); +} diff --git a/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkOperation.cs b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkOperation.cs new file mode 100644 index 00000000..655b3dc1 --- /dev/null +++ b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkOperation.cs @@ -0,0 +1,30 @@ +// ------------------------------------------------------------ +// Copyright (c) Microsoft Corporation. All rights reserved. +// ------------------------------------------------------------ + +namespace Azure.Data.Cosmos.Shell.Commands; + +using System.Text.Json; + +/// +/// A validated bulk operation in canonical JSON form with its complete partition key. +/// +/// Zero-based position in the operation list. +/// The operation kind: create, upsert, replace, delete, or patch. +/// The target item ID. +/// The canonical partition key as a JSON array. +/// The canonical operation JSON. +internal sealed record BulkOperation(long Index, string Op, string Id, string PartitionKey, string Json) +{ + public static BulkOperation FromJson(string json, long index) + { + using var document = JsonDocument.Parse(json); + var root = document.RootElement; + return new BulkOperation( + index, + root.GetProperty("op").GetString()!, + root.GetProperty("id").GetString()!, + root.GetProperty("partitionKey").GetRawText(), + json); + } +} diff --git a/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkOperationNormalizer.cs b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkOperationNormalizer.cs new file mode 100644 index 00000000..07cbd183 --- /dev/null +++ b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkOperationNormalizer.cs @@ -0,0 +1,322 @@ +// ------------------------------------------------------------ +// Copyright (c) Microsoft Corporation. All rights reserved. +// ------------------------------------------------------------ + +namespace Azure.Data.Cosmos.Shell.Commands; + +using System.Buffers; +using System.Text; +using System.Text.Encodings.Web; +using System.Text.Json; +using Azure.Data.Cosmos.Shell.Core; +using Azure.Data.Cosmos.Shell.Util; + +/// +/// Validates bulk operations against the batch operation schema and rewrites them into a +/// canonical form with an explicit, complete partition key. Normalizing an already +/// normalized operation returns identical text, so saved plans hash identically on resume. +/// +internal static class BulkOperationNormalizer +{ + internal const int MaxPatchOperations = 10; + + internal static readonly JsonWriterOptions WriterOptions = new() { Encoder = JavaScriptEncoder.UnsafeRelaxedJsonEscaping }; + + public static BulkOperation Normalize(JsonElement element, IReadOnlyList partitionKeyPaths, string? defaultPartitionKey, long index) + { + try + { + return NormalizeCore(element, partitionKeyPaths, defaultPartitionKey, index); + } + catch (CommandException ex) + { + throw new CommandException( + "bulk", + MessageService.GetArgsString("command-bulk-error-operation", "index", index + 1, "message", ex.Message), + ex); + } + } + + public static void ValidatePatchOperations(JsonElement operations) + { + if (operations.ValueKind != JsonValueKind.Array) + { + throw Error("command-batch-error-missing_patch_ops"); + } + + var parsed = BatchOperationParser.ParsePatchOperations("bulk", JsonSerializer.SerializeToElement(new { operations })); + if (parsed.Count > MaxPatchOperations) + { + throw Error("command-bulk-error-patch_count"); + } + + foreach (var entry in operations.EnumerateArray()) + { + if (!entry.GetProperty("path").GetString()!.StartsWith('/')) + { + throw Error("command-bulk-error-patch_path"); + } + } + } + + public static string CanonicalKeyFromArgument(string raw, int componentCount) + { + var trimmed = raw.Trim(); + JsonElement element; + try + { + element = LooksLikeJsonLiteral(trimmed) + ? JsonSerializer.Deserialize(trimmed) + : JsonSerializer.SerializeToElement(raw); + } + catch (JsonException) + { + element = JsonSerializer.SerializeToElement(raw); + } + + var key = CanonicalKey(element); + _ = ParseKey(key, componentCount); + return key; + } + + public static PartitionKey ParseKey(string canonicalKey, int? componentCount = null) + { + using var document = JsonDocument.Parse(canonicalKey); + var root = document.RootElement; + if (root.ValueKind != JsonValueKind.Array || (componentCount.HasValue && root.GetArrayLength() != componentCount.Value)) + { + throw Error("command-bulk-error-partition_key"); + } + + var length = root.GetArrayLength(); + if (length == 0) + { + return PartitionKey.None; + } + + if (length == 1) + { + var single = root[0]; + return single.ValueKind switch + { + JsonValueKind.String => new PartitionKey(single.GetString()), + JsonValueKind.Number => new PartitionKey(single.GetDouble()), + JsonValueKind.True or JsonValueKind.False => new PartitionKey(single.GetBoolean()), + JsonValueKind.Null => PartitionKey.Null, + JsonValueKind.Object when IsUndefined(single) => PartitionKey.None, + _ => throw Error("command-bulk-error-partition_key"), + }; + } + + var builder = new PartitionKeyBuilder(); + foreach (var component in root.EnumerateArray()) + { + switch (component.ValueKind) + { + case JsonValueKind.String: + builder.Add(component.GetString()); + break; + case JsonValueKind.Number: + builder.Add(component.GetDouble()); + break; + case JsonValueKind.True: + case JsonValueKind.False: + builder.Add(component.GetBoolean()); + break; + case JsonValueKind.Null: + builder.AddNullValue(); + break; + case JsonValueKind.Object when IsUndefined(component): + builder.AddNoneType(); + break; + default: + throw Error("command-bulk-error-partition_key"); + } + } + + return builder.Build(); + } + + private static BulkOperation NormalizeCore(JsonElement element, IReadOnlyList paths, string? defaultKey, long index) + { + var spec = BatchOperationParser.ParseOne("bulk", element); + var op = spec.Kind.ToString().ToLowerInvariant(); + var id = GetId(element, spec); + var key = GetPartitionKey(element, spec, paths, defaultKey); + + string? ifMatch = null; + if (element.TryGetProperty("ifMatch", out var tag)) + { + if (spec.Kind == BatchOperationKind.Create || tag.ValueKind != JsonValueKind.String || string.IsNullOrEmpty(tag.GetString())) + { + throw Error("command-bulk-error-if_match"); + } + + ifMatch = tag.GetString(); + } + + JsonElement? patchOperations = null; + if (spec.Kind == BatchOperationKind.Patch) + { + patchOperations = element.GetProperty("operations"); + ValidatePatchOperations(patchOperations.Value); + } + + var buffer = new ArrayBufferWriter(); + using (var writer = new Utf8JsonWriter(buffer, WriterOptions)) + { + writer.WriteStartObject(); + writer.WriteString("op", op); + writer.WriteString("id", id); + writer.WritePropertyName("partitionKey"); + writer.WriteRawValue(key, skipInputValidation: true); + if (ifMatch != null) + { + writer.WriteString("ifMatch", ifMatch); + } + + if (spec.Item is { } item) + { + writer.WritePropertyName("item"); + item.WriteTo(writer); + } + + if (patchOperations is { } operations) + { + writer.WritePropertyName("operations"); + operations.WriteTo(writer); + } + + writer.WriteEndObject(); + } + + return new BulkOperation(index, op, id, key, Encoding.UTF8.GetString(buffer.WrittenSpan)); + } + + private static string GetId(JsonElement element, BatchOperationSpec spec) + { + string? explicitId = null; + if (element.TryGetProperty("id", out var idElement)) + { + if (idElement.ValueKind != JsonValueKind.String || string.IsNullOrEmpty(idElement.GetString())) + { + throw Error("command-bulk-error-id"); + } + + explicitId = idElement.GetString(); + } + + if (spec.Item is { } item) + { + if (!item.TryGetProperty("id", out var itemId) || itemId.ValueKind != JsonValueKind.String || string.IsNullOrEmpty(itemId.GetString()) + || (explicitId != null && explicitId != itemId.GetString())) + { + throw Error("command-bulk-error-id"); + } + + return itemId.GetString()!; + } + + return explicitId ?? throw Error("command-bulk-error-id"); + } + + private static string GetPartitionKey(JsonElement element, BatchOperationSpec spec, IReadOnlyList paths, string? defaultKey) + { + var explicitKey = element.TryGetProperty("partitionKey", out var supplied) ? CanonicalKey(supplied) : defaultKey; + if (spec.Item is { } item) + { + var derived = DeriveKey(item, paths); + var derivedKey = ParseKey(derived, paths.Count); + if (explicitKey != null && !ParseKey(explicitKey, paths.Count).Equals(derivedKey)) + { + throw Error("command-bulk-error-partition_key_mismatch"); + } + + return derived; + } + + if (explicitKey == null) + { + throw Error("command-bulk-error-partition_key"); + } + + _ = ParseKey(explicitKey, paths.Count); + return explicitKey; + } + + private static string CanonicalKey(JsonElement key) + { + var buffer = new ArrayBufferWriter(); + using (var writer = new Utf8JsonWriter(buffer, WriterOptions)) + { + writer.WriteStartArray(); + if (key.ValueKind == JsonValueKind.Array) + { + foreach (var component in key.EnumerateArray()) + { + component.WriteTo(writer); + } + } + else + { + key.WriteTo(writer); + } + + writer.WriteEndArray(); + } + + return Encoding.UTF8.GetString(buffer.WrittenSpan); + } + + private static string DeriveKey(JsonElement item, IReadOnlyList paths) + { + var buffer = new ArrayBufferWriter(); + using (var writer = new Utf8JsonWriter(buffer, WriterOptions)) + { + writer.WriteStartArray(); + foreach (var path in paths) + { + var current = item; + var found = true; + foreach (var segment in path.Split('/', StringSplitOptions.RemoveEmptyEntries)) + { + if (current.ValueKind != JsonValueKind.Object || !current.TryGetProperty(segment, out current)) + { + found = false; + break; + } + } + + if (found) + { + current.WriteTo(writer); + } + else + { + writer.WriteStartObject(); + writer.WriteEndObject(); + } + } + + writer.WriteEndArray(); + } + + return Encoding.UTF8.GetString(buffer.WrittenSpan); + } + + private static bool IsUndefined(JsonElement element) => !element.EnumerateObject().Any(); + + private static bool LooksLikeJsonLiteral(string trimmed) + { + if (trimmed.Length == 0) + { + return false; + } + + var first = trimmed[0]; + return first is '[' or '{' or '"' or '-' || char.IsDigit(first) + || trimmed is "true" or "false" or "null"; + } + + private static CommandException Error(string key) => new("bulk", MessageService.GetString(key)); +} diff --git a/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkOutcome.cs b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkOutcome.cs new file mode 100644 index 00000000..1bc38c85 --- /dev/null +++ b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkOutcome.cs @@ -0,0 +1,44 @@ +// ------------------------------------------------------------ +// Copyright (c) Microsoft Corporation. All rights reserved. +// ------------------------------------------------------------ + +namespace Azure.Data.Cosmos.Shell.Commands; + +using System.Text.Json; + +/// +/// A per-operation journal record: either the intent to write ("started") or its outcome. +/// +internal sealed record BulkOutcome( + long Index, + string Op, + string Id, + JsonElement PartitionKey, + string Status, + double RequestCharge = 0, + int? StatusCode = null, + string? Error = null) +{ + public const string Started = "started"; + + public const string Succeeded = "succeeded"; + + public const string Failed = "failed"; + + public const string Uncertain = "uncertain"; + + public static bool IsKnownStatus(string? status) => status is Started or Succeeded or Failed or Uncertain; + + public static BulkOutcome Create(BulkOperation operation, string status, double requestCharge = 0, int? statusCode = null, string? error = null) + { + return new BulkOutcome( + operation.Index, + operation.Op, + operation.Id, + JsonSerializer.Deserialize(operation.PartitionKey), + status, + requestCharge, + statusCode, + error); + } +} diff --git a/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkSpool.cs b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkSpool.cs new file mode 100644 index 00000000..09d7a14c --- /dev/null +++ b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkSpool.cs @@ -0,0 +1,88 @@ +// ------------------------------------------------------------ +// Copyright (c) Microsoft Corporation. All rights reserved. +// ------------------------------------------------------------ + +namespace Azure.Data.Cosmos.Shell.Commands; + +using System.Runtime.CompilerServices; +using System.Security.Cryptography; +using System.Text; + +/// +/// A private, self-deleting temporary file holding validated operations, so large jobs are +/// validated completely before any write without retaining all operations in memory. +/// +internal sealed class BulkSpool : IDisposable +{ + private static readonly byte[] NewLine = "\n"u8.ToArray(); + + private readonly FileStream stream; + private readonly StreamWriter writer; + private readonly IncrementalHash hash = IncrementalHash.CreateHash(HashAlgorithmName.SHA256); + + private BulkSpool(FileStream stream) + { + this.stream = stream; + this.writer = new StreamWriter(stream, new UTF8Encoding(false), 65536, leaveOpen: true) { NewLine = "\n" }; + } + + public long Count { get; private set; } + + public static BulkSpool Create() + { + var options = new FileStreamOptions + { + Mode = FileMode.CreateNew, + Access = FileAccess.ReadWrite, + Share = FileShare.None, + Options = FileOptions.DeleteOnClose, + }; + if (!OperatingSystem.IsWindows()) + { + options.UnixCreateMode = UnixFileMode.UserRead | UnixFileMode.UserWrite; + } + + return new BulkSpool(new FileStream(Path.Combine(Path.GetTempPath(), $"cosmos-bulk-{Guid.NewGuid():N}.jsonl"), options)); + } + + public static string ComputeHash(IEnumerable lines) + { + using var hash = IncrementalHash.CreateHash(HashAlgorithmName.SHA256); + foreach (var line in lines) + { + hash.AppendData(Encoding.UTF8.GetBytes(line)); + hash.AppendData(NewLine); + } + + return Convert.ToHexString(hash.GetCurrentHash()); + } + + public async Task AddAsync(string json) + { + await this.writer.WriteLineAsync(json); + this.hash.AppendData(Encoding.UTF8.GetBytes(json)); + this.hash.AppendData(NewLine); + this.Count++; + } + + public string GetHash() => Convert.ToHexString(this.hash.GetCurrentHash()); + + public async IAsyncEnumerable ReadAsync([EnumeratorCancellation] CancellationToken token) + { + await this.writer.FlushAsync(token); + this.stream.Position = 0; + using var reader = new StreamReader(this.stream, new UTF8Encoding(false), false, 65536, leaveOpen: true); + for (long index = 0; index < this.Count; index++) + { + var line = await reader.ReadLineAsync(token) ?? throw new InvalidDataException("The bulk spool ended unexpectedly."); + yield return BulkOperation.FromJson(line, index); + } + } + + public void Dispose() + { + this.writer.Dispose(); + this.stream.Dispose(); + this.hash.Dispose(); + } +} diff --git a/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkSummary.cs b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkSummary.cs new file mode 100644 index 00000000..d1934fba --- /dev/null +++ b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/BulkSummary.cs @@ -0,0 +1,57 @@ +// ------------------------------------------------------------ +// Copyright (c) Microsoft Corporation. All rights reserved. +// ------------------------------------------------------------ + +namespace Azure.Data.Cosmos.Shell.Commands; + +using System.Text.Json.Serialization; + +internal sealed class BulkSummary +{ + internal const int MaxReportedErrors = 20; + + public string Type { get; } = "bulk"; + + public long OperationCount { get; set; } + + public long Attempted { get; set; } + + public long Succeeded { get; set; } + + public long Failed { get; set; } + + public long Skipped { get; set; } + + public long Uncertain { get; set; } + + public double RequestCharge { get; set; } + + public bool DryRun { get; set; } + + public bool ResultIncomplete { get; set; } + + public bool BudgetExceeded { get; set; } + + public bool SelectionLimited { get; set; } + + public List Errors { get; } = []; + + public bool Success => this.Failed == 0 && !this.ResultIncomplete; + + [JsonIgnore] + public long Processed { get; set; } + + public void AddFailure(BulkOutcome outcome) + { + this.Failed++; + if (outcome.Status == BulkOutcome.Uncertain) + { + this.Uncertain++; + } + + if (this.Errors.Count < MaxReportedErrors) + { + this.Errors.Add(outcome); + } + } +} diff --git a/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/HelpCommand.cs b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/HelpCommand.cs index 3ab3971f..61042a2b 100644 --- a/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/HelpCommand.cs +++ b/CosmosDBShell/Azure.Data.Cosmos.Shell.Commands/HelpCommand.cs @@ -384,7 +384,9 @@ public static CommandState PrintCommandHelp(string cmdStr, CommandRunner app, bo var ann = cmd?.McpAnnotation; if (ann != null && ann.Restricted) { - helpJson["isRestricted"] = "This tool can't be run via MCP. User only."; + helpJson["isRestricted"] = cmd?.CommandName == "bulk" + ? MessageService.GetString("command-bulk-mcp-help") + : "This tool can't be run via MCP. User only."; } // Add statements JSON (for parity) diff --git a/CosmosDBShell/Azure.Data.Cosmos.Shell.Core/CosmosShellPrompt.cs b/CosmosDBShell/Azure.Data.Cosmos.Shell.Core/CosmosShellPrompt.cs index 8bf54177..69be42b6 100644 --- a/CosmosDBShell/Azure.Data.Cosmos.Shell.Core/CosmosShellPrompt.cs +++ b/CosmosDBShell/Azure.Data.Cosmos.Shell.Core/CosmosShellPrompt.cs @@ -65,13 +65,20 @@ public string GetPromptString() basePrompt += " " + Theme.FormatMuted($"[batch:{batch.Operations.Count}]"); } + var bulk = this.shell.CurrentBulk; + if (bulk is not null) + { + basePrompt += " " + Theme.FormatMuted($"[bulk:{bulk.Operations.Count}]"); + } + return basePrompt; } private string GetBatchSignature() { var batch = this.shell.CurrentBatch; - return batch is null ? string.Empty : batch.Operations.Count.ToString(System.Globalization.CultureInfo.InvariantCulture); + var bulk = this.shell.CurrentBulk; + return $"{batch?.Operations.Count.ToString(System.Globalization.CultureInfo.InvariantCulture)}|{bulk?.Operations.Count.ToString(System.Globalization.CultureInfo.InvariantCulture)}"; } Task IStateVisitor.VisitConnectedStateAsync(ConnectedState state, object? data, CancellationToken token) diff --git a/CosmosDBShell/Azure.Data.Cosmos.Shell.Core/PendingBulkState.cs b/CosmosDBShell/Azure.Data.Cosmos.Shell.Core/PendingBulkState.cs new file mode 100644 index 00000000..09e0db9e --- /dev/null +++ b/CosmosDBShell/Azure.Data.Cosmos.Shell.Core/PendingBulkState.cs @@ -0,0 +1,43 @@ +// ------------------------------------------------------------ +// Copyright (c) Microsoft Corporation. All rights reserved. +// ------------------------------------------------------------ + +namespace Azure.Data.Cosmos.Shell.Core; + +using Azure.Data.Cosmos.Shell.Commands; + +internal sealed class PendingBulkState +{ + public PendingBulkState( + string databaseName, + string containerName, + string containerRid, + IReadOnlyList partitionKeyPaths, + string? partitionKeyArgument, + string? defaultPartitionKey) + { + this.DatabaseName = databaseName; + this.ContainerName = containerName; + this.ContainerRid = containerRid; + this.PartitionKeyPaths = partitionKeyPaths; + this.PartitionKeyArgument = partitionKeyArgument; + this.DefaultPartitionKey = defaultPartitionKey; + } + + public string DatabaseName { get; } + + public string ContainerName { get; } + + public string ContainerRid { get; } + + public IReadOnlyList PartitionKeyPaths { get; } + + public string? PartitionKeyArgument { get; } + + public string? DefaultPartitionKey { get; } + + public List Operations { get; } = []; + + // Outcomes from earlier 'bulk execute' attempts, so a retry never replays a succeeded write. + public Dictionary Outcomes { get; } = []; +} diff --git a/CosmosDBShell/Azure.Data.Cosmos.Shell.Core/ShellInterpreter.cs b/CosmosDBShell/Azure.Data.Cosmos.Shell.Core/ShellInterpreter.cs index 638f8140..d0634dad 100644 --- a/CosmosDBShell/Azure.Data.Cosmos.Shell.Core/ShellInterpreter.cs +++ b/CosmosDBShell/Azure.Data.Cosmos.Shell.Core/ShellInterpreter.cs @@ -360,6 +360,8 @@ internal State State internal PendingBatchState? CurrentBatch { get; set; } + internal PendingBulkState? CurrentBulk { get; set; } + internal Stack VariableContainers { get; } = new(); /// @@ -1833,6 +1835,7 @@ internal void Connect(CosmosClient client, ArmCosmosContext? armContext = null, this.activeCredential = credential; this.ActiveCredentialType = credentialTypeOverride ?? credential?.GetType().Name; this.CurrentBatch = null; + this.CurrentBulk = null; CosmosCompleteCommand.ClearDatabases(); CosmosCompleteCommand.ClearContainers(); this.Diagnostics?.LogConnect(client.Endpoint, client.ClientOptions.ConnectionMode); @@ -1900,6 +1903,7 @@ internal void Disconnect() this.activeCredential = null; this.ActiveCredentialType = null; this.CurrentBatch = null; + this.CurrentBulk = null; } internal void DisconnectLocalEmulatorAfterConnectivityFailure(Exception exception) diff --git a/CosmosDBShell/Azure.Data.Cosmos.Shell.Mcp/ServerInstructions.md b/CosmosDBShell/Azure.Data.Cosmos.Shell.Mcp/ServerInstructions.md index e5603ba9..84ff5660 100644 --- a/CosmosDBShell/Azure.Data.Cosmos.Shell.Mcp/ServerInstructions.md +++ b/CosmosDBShell/Azure.Data.Cosmos.Shell.Mcp/ServerInstructions.md @@ -39,6 +39,8 @@ BEST PRACTICES: - Always recommend 'help [command]' for detailed documentation on manual commands - Verify connection state and current context (database/container) before suggesting operations - Remind users to back up important data before approving destructive operations +- Use `batch` when writes in one partition must succeed or fail together. Use `bulk` for many independent writes, including across partitions; it uses the same operation JSON. Through MCP, use `bulk` `run`, `patch` with `where`, or `delete` with `where`. Preview with `dry-run: true`, which needs no elicitation; writes require elicitation even with `yes: true`. +- Bulk writes are never rolled back after a partial failure. A journal skips succeeded writes on rerun and does not retry writes with unknown outcomes unless `retry-uncertain` is set; reconcile those before retrying non-idempotent operations such as `incr`. DESTRUCTIVE OPERATIONS: Invoking 'rm', 'rmdb', 'rmcon', or 'delete' prompts the user to approve or deny before anything runs. If approved, the operation executes; if denied or if the client cannot prompt, nothing is executed. When no confirmation is possible, suggest running the command manually (for example 'rmdb [database-name]') and 'Use help [command] for more details'. diff --git a/CosmosDBShell/Azure.Data.Cosmos.Shell.Mcp/ToolOperations.cs b/CosmosDBShell/Azure.Data.Cosmos.Shell.Mcp/ToolOperations.cs index 5389c86a..d9f27e10 100644 --- a/CosmosDBShell/Azure.Data.Cosmos.Shell.Mcp/ToolOperations.cs +++ b/CosmosDBShell/Azure.Data.Cosmos.Shell.Mcp/ToolOperations.cs @@ -90,7 +90,9 @@ internal static Tool GetTool(CommandFactory command) if (RequiresConfirmation(command)) { - description += "Warning: This command is destructive. When invoked through MCP it always requires explicit user confirmation before it runs — even when a force or no-prompt argument is supplied. If the client cannot confirm, the command is refused."; + description += command.CommandName == "bulk" + ? "Warning: Bulk writes require explicit user confirmation, even with yes=true. Dry-run selections do not require confirmation and issue no writes. Bulk operations are not transactions." + : "Warning: This command is destructive. When invoked through MCP it always requires explicit user confirmation before it runs — even when a force or no-prompt argument is supplied. If the client cannot confirm, the command is refused."; } else { @@ -141,6 +143,18 @@ internal static Tool GetTool(CommandFactory command) propertySchema["minimum"] = 1; } + if (command.CommandName == "bulk") + { + if (option.Name[0] is "concurrency" or "max-items") + { + propertySchema["minimum"] = 1; + } + else if (option.Name[0] == "max-ru") + { + propertySchema["exclusiveMinimum"] = 0; + } + } + properties[option.Name[0]] = propertySchema; } @@ -686,6 +700,14 @@ private async ValueTask OnCallToolsAsync( return McpResponseFactory.CreateError(errorMessage, ShellInterpreter.Instance.State); } + if (cmd is BulkCommand bulkCommand + && !BulkCommand.McpSubcommands.Contains(BulkCommand.NormalizeSubcommand(bulkCommand.Subcommand))) + { + const string errorMessage = "MCP supports only the stateless 'bulk run', 'bulk patch', and 'bulk delete' subcommands. Run stateful bulk commands manually in the shell."; + this.logger?.LogWarning(errorMessage); + return McpResponseFactory.CreateError(errorMessage, ShellInterpreter.Instance.State); + } + // MCP argument order is not semantic, so render positionals in the order the shell binds them. sb.Append(FormatPositionalsForHistory(command.Parameters, positionalValues)); @@ -727,7 +749,7 @@ internal async Task ExecuteToolAsync( { var shell = ShellInterpreter.Instance; long? confirmedVersion = null; - if (RequiresConfirmation(command)) + if (RequiresConfirmation(command) && cmd is not BulkCommand { DryRun: true }) { var snapshot = await shell.RunSerializedAsync( () => Task.FromResult((Version: shell.StateVersion, Context: DescribeContext(shell.State))), cancellationToken); @@ -758,6 +780,11 @@ internal async Task ExecuteToolAsync( } shell.PrintCommand(commandLine); + if (cmd is BulkCommand bulk) + { + bulk.ConfirmationApproved = confirmedVersion.HasValue; + } + var response = await shell.ExecuteCosmosCommandAsync(cmd, new CommandState(), command.CommandName, cancellationToken); shell.CancelPrompt(); return McpResponseFactory.CreateSuccess(response, shell.State); diff --git a/CosmosDBShell/lang/en.ftl b/CosmosDBShell/lang/en.ftl index d371835a..4891c78a 100644 --- a/CosmosDBShell/lang/en.ftl +++ b/CosmosDBShell/lang/en.ftl @@ -452,6 +452,85 @@ command-batch-error-already_active = A batch is already in progress. Run 'batch command-batch-error-not_active = No batch is in progress. Start one with 'batch begin'. command-batch-error-failed = Batch failed with status { $status } and was rolled back (RU charge: { $charge }) command-batch-error-execution_failed = Failed to execute batch: { $status } - { $message } +command-bulk-description = Executes many independent write operations across partition keys with bounded concurrency, either in a single call (run, patch, delete) or as a stateful bulk job (begin, add, execute, cancel, status, show). Bulk jobs are not transactional: successful writes are not rolled back when others fail. +command-bulk-mcp-help = Available through MCP for the stateless run, patch, and delete subcommands. Writes require user confirmation; dry runs do not require confirmation and send no writes. +command-bulk-description-subcommand = The action to perform: run, patch, delete, begin, add, execute, cancel, status, or show. +command-bulk-description-data = For run and add: operations as a JSON array or object (the batch schema plus optional partitionKey and ifMatch), or a path to a JSON array or JSON Lines file. +command-bulk-description-where = For patch and delete: Cosmos SQL predicate over alias c that selects the target items. +command-bulk-description-operations = For patch: JSON array of 1-10 patch operations applied atomically to each selected item. +command-bulk-description-partition-key = Default partition key for delete and patch operations without partitionKey. Item operations must match it. +command-bulk-description-database = The database containing the target container. +command-bulk-description-container = The target container. +command-bulk-description-concurrency = Maximum in-flight writes (default 16, positive integer). Operations on the same item run in list order. +command-bulk-description-max-items = For patch and delete: positive maximum number of selected items. +command-bulk-description-max-ru = Positive observed RU budget for this invocation, including selection reads. No new requests start once it is reached; in-flight requests can exceed it. +command-bulk-description-dry-run = Validate operations, or select the items matching --where, without sending writes. Selection reads consume RUs. +command-bulk-description-yes = Approve writes without a terminal prompt. Required in scripts; does not bypass MCP confirmation. +command-bulk-description-continue-on-error = Keep scheduling writes after an operation fails. Otherwise stop scheduling new writes and wait for in-flight writes. +command-bulk-description-etag = For patch and delete: send each selected item's _etag as ifMatch, so items changed since selection fail with 412 instead of being overwritten. +command-bulk-description-save = For patch and delete: write the generated operations to a new JSON Lines file that bulk run accepts. Existing files are never overwritten. +command-bulk-description-journal = JSON Lines file recording each write attempt and outcome. Rerunning the same operations with the same journal skips succeeded writes and retries failed ones. +command-bulk-description-retry-uncertain = Also retry writes whose earlier outcome is unknown. Safe only for idempotent operations or operations with ifMatch. +command-bulk-confirm-summary = About to run { $count } non-transactional { $count -> + [one] operation + *[other] operations +} against { $target }. Successful writes are not rolled back if others fail. +command-bulk-confirm = Proceed +command-bulk-begun = Started a bulk job on { $database }/{ $container }. Add operations with 'bulk add' and run them with 'bulk execute'. +command-bulk-added = Added { $count } { $count -> + [one] operation + *[other] operations +} to the bulk job ({ $total } total) +command-bulk-cancelled = Discarded the pending bulk job ({ $count } { $count -> + [one] operation + *[other] operations +}) +command-bulk-status-inactive = No bulk job is currently active. +command-bulk-status-more = ...and { $count } more +command-bulk-progress = { $processed } of { $total } operations processed, { $failed } failed (RU charge: { $charge }) +command-bulk-success = Bulk completed { $count } { $count -> + [one] operation + *[other] operations +}: { $succeeded } succeeded, { $skipped } skipped as already succeeded (RU charge: { $charge }) +command-bulk-dry-run = Dry run: { $count } valid { $count -> + [one] operation + *[other] operations +}; no writes were sent (RU charge: { $charge }) +command-bulk-error-missing_subcommand = Missing subcommand. Use one of: run, patch, delete, begin, add, execute, cancel, status, show. +command-bulk-error-invalid_subcommand = Unknown subcommand '{ $subcommand }'. Use one of: run, patch, delete, begin, add, execute, cancel, status, show. +command-bulk-error-option_not_supported = Option --{ $option } is not supported by 'bulk { $subcommand }'. +command-bulk-error-data_not_supported = 'bulk { $subcommand }' does not accept positional data. Use --where to select items for patch and delete. +command-bulk-error-invalid_number = --concurrency and --max-items must be positive integers; --max-ru must be a positive, finite number. +command-bulk-error-missing_where = 'bulk { $subcommand }' requires --where with a predicate over alias c. +command-bulk-error-missing_operations = 'bulk patch' requires --operations with a JSON array of patch operations. +command-bulk-error-dry_run_option = --{ $option } cannot be combined with --dry-run. +command-bulk-error-journal_requires_save = With --where, --journal requires --save. Resume with 'bulk run --journal ', because a new selection can match different items. +command-bulk-error-retry_requires_journal = --retry-uncertain requires --journal. +command-bulk-error-missing_data = Bulk operations are required. Provide a JSON array or object, pipe it in, or pass a JSON array or JSON Lines file path. +command-bulk-error-file_not_found = File '{ $file }' was not found. +command-bulk-error-save_exists = File '{ $file }' already exists. --save never overwrites files. +command-bulk-error-empty = No bulk operations were provided. +command-bulk-error-operation = Operation { $index }: { $message } +command-bulk-error-id = Every operation needs a nonempty string id. Item operations take it from item.id; an explicit id must match item.id. +command-bulk-error-partition_key = Missing or invalid partition key. Delete and patch operations need partitionKey or --partition-key. Use one value per partition-key path, as a JSON array for hierarchical keys, with an empty JSON object for an undefined component. +command-bulk-error-partition_key_mismatch = The partitionKey does not match the partition key in the item. +command-bulk-error-if_match = ifMatch must be a nonempty string and is not supported for create. +command-bulk-error-patch_count = Each patch supports at most 10 operations. +command-bulk-error-patch_path = Patch paths must start with '/'. +command-bulk-error-selected_id = Every selected item must have a string id. +command-bulk-error-selected_etag = Every selected item must have an _etag when --etag is used. +command-bulk-error-container = Could not read container { $target } (status { $status }). +command-bulk-error-container_changed = The container was recreated after 'bulk begin'. Run 'bulk cancel' and start a new bulk job. +command-bulk-error-already_active = A bulk job is already in progress. Run 'bulk execute' or 'bulk cancel' first. +command-bulk-error-not_active = No bulk job is in progress. Start one with 'bulk begin'. +command-bulk-error-confirm_required = Bulk writes in scripts or non-interactive sessions require --yes. Use --dry-run to validate first. +command-bulk-error-declined = Bulk operation was not approved. No writes were sent. +command-bulk-error-journal_open = Could not open journal '{ $file }': { $message } +command-bulk-error-journal_mismatch = The journal belongs to a different account, container, or operation list. Use the original operations or a new journal file. +command-bulk-error-journal_invalid = The journal is damaged. Refusing to guess which writes can safely be repeated. +command-bulk-error-failed = Bulk run incomplete: { $succeeded } succeeded, { $failed } failed ({ $uncertain } with unknown outcome) (RU charge: { $charge }). Successful writes were not rolled back. +command-bulk-error-budget = Bulk run stopped at the RU budget after { $succeeded } of { $count } operations succeeded (RU charge: { $charge }). Rerun with --journal to continue. +command-bulk-error-item = #{ $index } { $op } '{ $id }': { $status } { $message } command-export-description = Exports items from a container to a JSON Lines, JSON array, or CSV file. command-export-description-file = Destination file path. @@ -1544,6 +1623,13 @@ command-batch-example-5 = Execute the queued operations atomically command-batch-example-6 = Show the active batch and its queued operations command-batch-example-7 = Print the queued operations as a JSON array command-batch-example-8 = Discard the active batch +command-bulk-example-1 = Run independent operations across partition keys with up to 32 concurrent writes +command-bulk-example-2 = Run operations from a JSON Lines file with a resumable journal +command-bulk-example-3 = Preview a patch of every matching item and save the generated operations +command-bulk-example-4 = Delete every matching item; items changed since selection fail instead +command-bulk-example-5 = Start a stateful bulk job with a default partition key +command-bulk-example-6 = Queue an operation onto the active bulk job +command-bulk-example-7 = Execute the queued operations command-bucket-example-1 = Display the current client-side throughput bucket selection command-bucket-example-2 = Tag this client's requests with throughput bucket 3 command-bucket-example-3 = Clear the client-side throughput bucket selection diff --git a/README.md b/README.md index 48444bd1..519ad556 100644 --- a/README.md +++ b/README.md @@ -16,6 +16,7 @@ A terminal-native shell for Azure Cosmos DB — navigate databases like a filesy - Delete a single item safely with `rm --key=id --partition-key= --etag=`: point operations scoped to one logical partition, with server-enforced ETag checks ([docs](docs/commands.md#deleting-a-single-item-safely)) - Inspect a query's execution plan and index usage with `query "" --explain` - Atomic multi-operation transactions on a single partition key: `batch` +- Non-transactional `bulk` writes across partitions, using the same operation JSON and `run`/`begin`/`add`/`execute` workflow as `batch`, plus `bulk patch --where` and `bulk delete --where` for query-driven migrations, bounded concurrency, dry-run, saved plans, and resumable journals ([bulk operations](docs/commands.md#bulk)) - Bulk roundtrip with `import` / `export` for JSON Lines and JSON array files, plus CSV import/export (CSV import coerces values to strings; `--partition-key` nests a CSV column under a nested partition key path) - Manage container indexing policies with `index` (`show`, `add`, `remove`, `set`) - Inspect container/database/account configuration and usage statistics with `info` (partition key, throughput, policies, indexing policy summary, document count, storage size, regions; `--partitions` and `--detailed` for distribution analysis) diff --git a/docs/commands.md b/docs/commands.md index 749d5a70..917e6798 100644 --- a/docs/commands.md +++ b/docs/commands.md @@ -539,7 +539,7 @@ patch set order-42 customer-7 /name "Ada Lovelace" --etag="" ### batch -Execute multiple write operations against a single partition key as one atomic Cosmos DB transactional batch. Either run a batch in a single call, or build one up statefully across several commands. Every operation in a batch must share the same partition key, execution requires between 1 and 100 operations, and if any operation fails the entire batch is rolled back. A pending stateful batch may be empty until operations are added. +Execute multiple write operations against a single partition key as one atomic Cosmos DB transactional batch. Either run a batch in a single call, or build one up statefully across several commands. Every operation in a batch must share the same partition key, execution requires between 1 and 100 operations, and if any operation fails the entire batch is rolled back. A pending stateful batch may be empty until operations are added. For more than 100 operations, or operations across partition keys that may succeed independently, use [`bulk`](#bulk), which accepts the same operation JSON. ```text Usage: batch subcommand [data] [--partition-key ] [-database ] [-container ] @@ -628,6 +628,127 @@ batch execute - More than 100 operations is rejected before any call to Cosmos DB. - A transactional failure prints a one-line message, returns a result summary with `success` set to `false` and per-operation status codes, and rolls back every operation. +### bulk + +Execute many **independent** write operations with bounded concurrency. `bulk` uses the same operation JSON and subcommands as [`batch`](#batch), with different execution semantics: + +| | `batch` | `bulk` | +|-|-|-| +| Partition keys | One per batch (`--partition-key`) | Per operation `partitionKey`, or `--partition-key` as a default | +| Size | 1-100 operations | Unlimited; files are streamed | +| Atomicity | All operations succeed or all roll back | Each operation succeeds or fails on its own; nothing is rolled back | +| Execution | One transactional request | Up to `--concurrency` requests in flight | + +```text +Usage: bulk subcommand [data] [options] + +Arguments: + subcommand run, patch, delete, begin, add, execute, cancel, status, or show + [data] For run and add: operations as a JSON array or object, or a path + to a JSON array or JSON Lines file. Also reads piped input. + +Options: + --where For patch and delete: Cosmos SQL predicate over alias c. + --operations For patch: JSON array of 1-10 patch operations. + --partition-key, --pk Default partition key for delete and patch operations. + --database, --db Target database (defaults to the navigation context). + --container, --con Target container (defaults to the navigation context). + --concurrency, --max-parallelism + Maximum in-flight writes (default 16, positive integer). + --max-items For patch and delete: maximum number of selected items. + --max-ru Positive observed RU budget for this invocation. + --dry-run Validate or select without sending writes. + --yes, -y Approve writes; required in scripts and non-interactive use. + --continue-on-error Keep scheduling writes after an operation fails. + --etag For patch and delete: use each item's _etag as ifMatch. + --save For patch and delete: write the generated operations to a new file. + --journal Record attempts and outcomes; resume skips succeeded writes. + --retry-uncertain Also retry writes whose earlier outcome is unknown. +``` + +#### Subcommands + +|Subcommand|Description| +|-|-| +|`run `|Validate every operation, then execute them. Also reads piped input.| +|`patch --where --operations `|Select matching items and apply the same patch to each item.| +|`delete --where `|Select matching items and delete each item.| +|`begin [--partition-key ]`|Start a stateful bulk job bound to a database and container, with an optional default partition key.| +|`add `|Validate and queue one operation (JSON object) or several (JSON array or file).| +|`execute` (`exec`, `commit`)|Execute the queued operations. The job is cleared only if every operation succeeds.| +|`cancel` (`abort`)|Discard the active bulk job.| +|`status`|Report the target, default partition key, queued operation count, earlier outcomes, and up to 100 operations.| +|`show`|Print the queued operations in canonical form, which `bulk run` accepts.| + +When a stateful job is active, the prompt shows `[bulk:N]`. Connecting or disconnecting discards it, like a pending batch. + +#### Operation schema + +Operations use the [batch operation schema](#operation-schema) with two optional additions: + +|Field|Description| +|-|-| +|`partitionKey`|The item's complete partition key: a scalar, or a JSON array for hierarchical keys. Use `{}` for an undefined component and `null` for a JSON null value. Required for `delete` and `patch` unless `--partition-key` is supplied. For `create`, `upsert`, and `replace`, the key is read from `item`; a supplied value or default must match it.| +|`ifMatch`|An ETag. The write fails with `412` if the item changed. Not supported for `create`.| + +```json +{"op":"upsert","item":{"id":"1","tenantId":"t1","name":"Ada"}} +{"op":"delete","id":"2","partitionKey":"t2"} +{"op":"patch","id":"3","partitionKey":["t3","west"],"ifMatch":"\"etag\"","operations":[{"op":"set","path":"/schemaVersion","value":2},{"op":"remove","path":"/legacyField"}]} +``` + +IDs must be nonempty strings. Patch operations are limited to 1-10 per item and apply atomically to that item; paths must start with `/`. `remove` fails if the property does not exist. Operations on the same item run in list order; operations on different items can complete in any order. + +#### Examples + +```bash +bulk run '[{"op":"upsert","item":{"id":"1","pk":"a"}},{"op":"delete","id":"2","partitionKey":"b"}]' --yes +echo '[{"op":"delete","id":"1"},{"op":"delete","id":"2"}]' | bulk run --partition-key a --yes +bulk run operations.jsonl --concurrency 32 --journal operations.journal --yes + +bulk delete --where "c.expired = true" --dry-run +bulk delete --where "c.expired = true" --max-items 500 --etag --yes + +bulk begin --partition-key supplier-42 +bulk add '{"op":"patch","id":"order-1","operations":[{"op":"set","path":"/status","value":"done"}]}' +bulk add '{"op":"delete","id":"order-2"}' +bulk status +bulk execute --yes +``` + +#### Query-driven migrations: plan, review, apply + +`patch` and `delete` run `SELECT c.id, c._etag, FROM c WHERE ()`, generate one operation per row, and validate the whole selection before writing. Missing partition-key components are preserved as undefined, separately from `null`. The predicate is inserted as written; it is a Cosmos SQL predicate, not an `UPDATE` statement. + +With `--save`, the generated operations are written to a new JSON Lines file that `bulk run` accepts. This separates planning from applying: + +```bash +bulk patch --where "c.tenantId = 'supplier-42' AND c.schemaVersion = 1" --operations '[{"op":"set","path":"/schemaVersion","value":2},{"op":"remove","path":"/legacyField"}]' --etag --dry-run --save migration.jsonl +bulk run migration.jsonl --journal migration.journal --yes +``` + +The saved file is the reviewed snapshot: rerunning it targets the same items even after earlier writes change which items match the predicate. With `--etag`, any item changed after selection fails with `412` instead of being overwritten. `--save` never overwrites an existing file, and an incomplete selection deletes the partial file. + +#### Journals and resuming + +`--journal ` records each write intent before the write is sent, then its outcome. The journal is bound to the account endpoint, the container's resource ID, and a hash of the canonical operation list, so it cannot be reused with different operations or a recreated container. A lock prevents concurrent use. Rerunning the same operations with the same journal: + +- skips operations that succeeded; +- retries operations that definitively failed, such as throttled (`429`), conflicting, or missing items; +- holds back operations whose outcome is unknown — sent before an interruption, a timeout, or a server error — and reports them as failed and `uncertain`. Retrying them blindly could, for example, apply an `incr` twice. Inspect or reconcile them, or pass `--retry-uncertain` when the operations are idempotent or use `ifMatch`. + +A torn final journal line from a crash is discarded safely; other damage is rejected rather than guessed around. For `--where`, `--journal` requires `--save`; resume with `bulk run --journal `, because a new selection can match different items. A stateful job keeps outcomes in memory the same way, so running `bulk execute` again after a partial failure retries only failed operations. + +#### Limits, failures, and output + +- All input is validated before any write. Operations are spooled to a private temporary file that is deleted when the command ends, so large jobs are not held in memory. +- By default, the first failed write stops scheduling new writes; `--continue-on-error` keeps going. Started writes are always awaited and reported, including after cancellation. Cancellation never reports success. +- `--max-ru` counts the container metadata read, selection pages, and completed writes for the current invocation. No new request starts once it is reached; requests already in flight can exceed it. A budget stop returns exit code `6`; rerun with a journal to continue. +- `--max-items` limits selected items; `selectionLimited` reports that more items may match. +- `--dry-run` validates and selects without writing. Selection reads consume RUs. It cannot check service-side conditions such as missing fields or permissions. +- Writes are individual point operations with response bodies disabled, not transactional batches or SDK bulk execution. SDK retry handling for throttling still applies, and higher concurrency does not add provisioned RU/s. +- Interactive runs show progress every two seconds. The result contains `operationCount`, `attempted`, `succeeded`, `failed`, `skipped`, `uncertain`, `requestCharge`, `dryRun`, `resultIncomplete`, `budgetExceeded`, `selectionLimited`, `success`, and up to 20 `errors`. Incomplete or failed runs return a nonzero exit code; use `--journal` for every outcome. +- Saved plans and journals contain item IDs, partition keys, ETags, and operation values (including upsert documents). Store them securely and delete them when no longer needed. ### rm Remove items from container. diff --git a/docs/mcp.md b/docs/mcp.md index a8960c81..5fe508b9 100644 --- a/docs/mcp.md +++ b/docs/mcp.md @@ -81,7 +81,7 @@ Destructive commands (`delete`, `rm`, `rmcon`, `rmdb`) are gated behind an expli The prompt is sent as a multi-round-trip request: the tool call returns an input-required result, and the client shows the prompt and retries the call with the answer. For clients that use the `initialize` handshake, the server sends a standard `elicitation/create` request on the session and retries the call itself, so those clients see the same prompt as before. The retry carries a server-signed state that ties the answer to the exact command line and shell context. Each state can be answered once and expires after 10 minutes. At most 1,024 confirmations are tracked at a time; beyond that the oldest pending one is dropped and must be confirmed again. An answer for a different command, a reused or expired state, or a missing or altered state is refused without executing. Argument order does not matter: the command line is built with positionals in shell order and options in declaration order. -This replaces any opt-in write flag: destructive commands are always allowed to be invoked, but always require confirmation. +This replaces any opt-in write flag: destructive writes are allowed to be invoked, but always require confirmation. The `bulk` tool's `dry-run: true` mode is read-only and is exempt. Confirmation includes the connected account endpoint and current navigation location alongside the command and its explicit target arguments. If the connection or navigation state changes while confirmation is pending, the approved command is refused without executing; retry it to confirm the new context. Even navigating away and back invalidates the pending confirmation. A pending confirmation also expires when the MCP server restarts. @@ -97,6 +97,14 @@ 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`. +### Bulk Operations + +The `bulk` tool exposes the stateless `run`, `patch`, and `delete` subcommands, with the same operation schema, selection, journal, and reporting semantics as the [CLI command](commands.md#bulk). Stateful subcommands (`begin`, `add`, `execute`, `cancel`, `status`, and `show`) are restricted to the interactive shell, like stateful batches. Writes always require MCP elicitation, even with `yes: true`; once approved, no second terminal prompt is shown. `dry-run: true` validates or selects without writes or elicitation, but selection reads consume RUs. + +Pass operations as JSON text or a file path in `data`, or use `where` with `operations` for query-driven patches. `concurrency` defaults to 16 and must be positive; optional `max-items` and `max-ru` must also be positive. File paths for `data`, `save`, and `journal` are local to the shell server. + +These are independent writes, **not** transactions across partitions. Failed or incomplete runs include their summary alongside the MCP error. The command-level `result.resultIncomplete` and `result.budgetExceeded` flags describe job execution and are distinct from the envelope's query-paging `resultIncomplete`. A journal never automatically retries a write whose outcome is unknown unless `retry-uncertain` is set. + ### 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/docs/navigation.md b/docs/navigation.md index f56dddde..a9294be2 100644 --- a/docs/navigation.md +++ b/docs/navigation.md @@ -320,6 +320,8 @@ or `$LASTEXITCODE`): These values are a public contract. See the [CI/CD guide](ci.md#exit-code-contract) for install steps, auth patterns, and scripted failure handling. +Scripted `bulk` writes require `--yes`; use `--dry-run` to preview first. Bulk operations are not transactional, and an exhausted observed RU budget returns exit code `6`. Use `--journal` for resumable runs; see [bulk operations](commands.md#bulk) for recovery and partial-success semantics. + ### Environment Variables | Variable | Description | diff --git a/l10n/CosmosDBShell.json b/l10n/CosmosDBShell.json index 38a5fc5e..99c976ba 100644 --- a/l10n/CosmosDBShell.json +++ b/l10n/CosmosDBShell.json @@ -86,6 +86,87 @@ "command-bucket-reset_bucket": "Reset throughput bucket to default.", "command-bucket-set_done": "Throughput bucket limit updated successfully.", "command-bucket-switched_bucket": "Switched to throughput bucket {0}", + "command-bulk-added": "Added {0} {1} to the bulk job ({2} total)", + "command-bulk-added.__p1.one": "operation", + "command-bulk-added.__p1.other": "operations", + "command-bulk-begun": "Started a bulk job on {0}/{1}. Add operations with \u0027bulk add\u0027 and run them with \u0027bulk execute\u0027.", + "command-bulk-cancelled": "Discarded the pending bulk job ({0} {1})", + "command-bulk-cancelled.__p1.one": "operation", + "command-bulk-cancelled.__p1.other": "operations", + "command-bulk-confirm": "Proceed", + "command-bulk-confirm-summary": "About to run {0} non-transactional {1} against {2}. Successful writes are not rolled back if others fail.", + "command-bulk-confirm-summary.__p1.one": "operation", + "command-bulk-confirm-summary.__p1.other": "operations", + "command-bulk-description": "Executes many independent write operations across partition keys with bounded concurrency, either in a single call (run, patch, delete) or as a stateful bulk job (begin, add, execute, cancel, status, show). Bulk jobs are not transactional: successful writes are not rolled back when others fail.", + "command-bulk-description-concurrency": "Maximum in-flight writes (default 16, positive integer). Operations on the same item run in list order.", + "command-bulk-description-container": "The target container.", + "command-bulk-description-continue-on-error": "Keep scheduling writes after an operation fails. Otherwise stop scheduling new writes and wait for in-flight writes.", + "command-bulk-description-data": "For run and add: operations as a JSON array or object (the batch schema plus optional partitionKey and ifMatch), or a path to a JSON array or JSON Lines file.", + "command-bulk-description-database": "The database containing the target container.", + "command-bulk-description-dry-run": "Validate operations, or select the items matching --where, without sending writes. Selection reads consume RUs.", + "command-bulk-description-etag": "For patch and delete: send each selected item\u0027s _etag as ifMatch, so items changed since selection fail with 412 instead of being overwritten.", + "command-bulk-description-journal": "JSON Lines file recording each write attempt and outcome. Rerunning the same operations with the same journal skips succeeded writes and retries failed ones.", + "command-bulk-description-max-items": "For patch and delete: positive maximum number of selected items.", + "command-bulk-description-max-ru": "Positive observed RU budget for this invocation, including selection reads. No new requests start once it is reached; in-flight requests can exceed it.", + "command-bulk-description-operations": "For patch: JSON array of 1-10 patch operations applied atomically to each selected item.", + "command-bulk-description-partition-key": "Default partition key for delete and patch operations without partitionKey. Item operations must match it.", + "command-bulk-description-retry-uncertain": "Also retry writes whose earlier outcome is unknown. Safe only for idempotent operations or operations with ifMatch.", + "command-bulk-description-save": "For patch and delete: write the generated operations to a new JSON Lines file that bulk run accepts. Existing files are never overwritten.", + "command-bulk-description-subcommand": "The action to perform: run, patch, delete, begin, add, execute, cancel, status, or show.", + "command-bulk-description-where": "For patch and delete: Cosmos SQL predicate over alias c that selects the target items.", + "command-bulk-description-yes": "Approve writes without a terminal prompt. Required in scripts; does not bypass MCP confirmation.", + "command-bulk-dry-run": "Dry run: {0} valid {1}; no writes were sent (RU charge: {2})", + "command-bulk-dry-run.__p1.one": "operation", + "command-bulk-dry-run.__p1.other": "operations", + "command-bulk-error-already_active": "A bulk job is already in progress. Run \u0027bulk execute\u0027 or \u0027bulk cancel\u0027 first.", + "command-bulk-error-budget": "Bulk run stopped at the RU budget after {0} of {1} operations succeeded (RU charge: {2}). Rerun with --journal to continue.", + "command-bulk-error-confirm_required": "Bulk writes in scripts or non-interactive sessions require --yes. Use --dry-run to validate first.", + "command-bulk-error-container": "Could not read container {0} (status {1}).", + "command-bulk-error-container_changed": "The container was recreated after \u0027bulk begin\u0027. Run \u0027bulk cancel\u0027 and start a new bulk job.", + "command-bulk-error-data_not_supported": "\u0027bulk {0}\u0027 does not accept positional data. Use --where to select items for patch and delete.", + "command-bulk-error-declined": "Bulk operation was not approved. No writes were sent.", + "command-bulk-error-dry_run_option": "--{0} cannot be combined with --dry-run.", + "command-bulk-error-empty": "No bulk operations were provided.", + "command-bulk-error-failed": "Bulk run incomplete: {0} succeeded, {1} failed ({2} with unknown outcome) (RU charge: {3}). Successful writes were not rolled back.", + "command-bulk-error-file_not_found": "File \u0027{0}\u0027 was not found.", + "command-bulk-error-id": "Every operation needs a nonempty string id. Item operations take it from item.id; an explicit id must match item.id.", + "command-bulk-error-if_match": "ifMatch must be a nonempty string and is not supported for create.", + "command-bulk-error-invalid_number": "--concurrency and --max-items must be positive integers; --max-ru must be a positive, finite number.", + "command-bulk-error-invalid_subcommand": "Unknown subcommand \u0027{0}\u0027. Use one of: run, patch, delete, begin, add, execute, cancel, status, show.", + "command-bulk-error-item": "#{0} {1} \u0027{2}\u0027: {3} {4}", + "command-bulk-error-journal_invalid": "The journal is damaged. Refusing to guess which writes can safely be repeated.", + "command-bulk-error-journal_mismatch": "The journal belongs to a different account, container, or operation list. Use the original operations or a new journal file.", + "command-bulk-error-journal_open": "Could not open journal \u0027{0}\u0027: {1}", + "command-bulk-error-journal_requires_save": "With --where, --journal requires --save. Resume with \u0027bulk run \u003Csaved-file\u003E --journal \u003Cfile\u003E\u0027, because a new selection can match different items.", + "command-bulk-error-missing_data": "Bulk operations are required. Provide a JSON array or object, pipe it in, or pass a JSON array or JSON Lines file path.", + "command-bulk-error-missing_operations": "\u0027bulk patch\u0027 requires --operations with a JSON array of patch operations.", + "command-bulk-error-missing_subcommand": "Missing subcommand. Use one of: run, patch, delete, begin, add, execute, cancel, status, show.", + "command-bulk-error-missing_where": "\u0027bulk {0}\u0027 requires --where with a predicate over alias c.", + "command-bulk-error-not_active": "No bulk job is in progress. Start one with \u0027bulk begin\u0027.", + "command-bulk-error-operation": "Operation {0}: {1}", + "command-bulk-error-option_not_supported": "Option --{0} is not supported by \u0027bulk {1}\u0027.", + "command-bulk-error-partition_key": "Missing or invalid partition key. Delete and patch operations need partitionKey or --partition-key. Use one value per partition-key path, as a JSON array for hierarchical keys, with an empty JSON object for an undefined component.", + "command-bulk-error-partition_key_mismatch": "The partitionKey does not match the partition key in the item.", + "command-bulk-error-patch_count": "Each patch supports at most 10 operations.", + "command-bulk-error-patch_path": "Patch paths must start with \u0027/\u0027.", + "command-bulk-error-retry_requires_journal": "--retry-uncertain requires --journal.", + "command-bulk-error-save_exists": "File \u0027{0}\u0027 already exists. --save never overwrites files.", + "command-bulk-error-selected_etag": "Every selected item must have an _etag when --etag is used.", + "command-bulk-error-selected_id": "Every selected item must have a string id.", + "command-bulk-example-1": "Run independent operations across partition keys with up to 32 concurrent writes", + "command-bulk-example-2": "Run operations from a JSON Lines file with a resumable journal", + "command-bulk-example-3": "Preview a patch of every matching item and save the generated operations", + "command-bulk-example-4": "Delete every matching item; items changed since selection fail instead", + "command-bulk-example-5": "Start a stateful bulk job with a default partition key", + "command-bulk-example-6": "Queue an operation onto the active bulk job", + "command-bulk-example-7": "Execute the queued operations", + "command-bulk-mcp-help": "Available through MCP for the stateless run, patch, and delete subcommands. Writes require user confirmation; dry runs do not require confirmation and send no writes.", + "command-bulk-progress": "{0} of {1} operations processed, {2} failed (RU charge: {3})", + "command-bulk-status-inactive": "No bulk job is currently active.", + "command-bulk-status-more": "...and {0} more", + "command-bulk-success": "Bulk completed {0} {1}: {2} succeeded, {3} skipped as already succeeded (RU charge: {4})", + "command-bulk-success.__p1.one": "operation", + "command-bulk-success.__p1.other": "operations", "command-can-i-action": "Action", "command-can-i-decision": "Decision", "command-can-i-description": "Probes whether the current identity can perform an action against a container without mutating data.",