Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 17 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,10 +6,27 @@

- Cosmos DB data-plane commands now consistently expose their aggregate observed request charge in structured output and connection-scoped `info` telemetry, including metadata/configuration operations, scripts, change feed reads, paginated operations, handled probes, and charged failures. Azure Resource Manager control-plane operations remain uncharged.
- Added `$sessionRequestCharge` and `$sessionChargedOperationCount` as read-only shell variables. Set `$sessionRequestChargeWarningThreshold` to a positive RU threshold to print one warning when the current connection reaches it; `info` reports it as `session.requestChargeWarningThreshold`.
- Destructive MCP confirmations now identify their target. The elicitation prompt adds the connected account endpoint and the current database/container location, and notes that explicit `--db`/`--con` arguments override that location. ([#207](https://github.com/Azure/CosmosDBShell/pull/207))
- Import and export no longer hold entire files in memory. CSV imports are parsed incrementally, and CSV exports spool documents to a private temporary file to determine the complete column set, so transfers no longer scale with document count. Allow temporary disk space for the CSV export spool in addition to the destination file. ([#207](https://github.com/Azure/CosmosDBShell/pull/207))

### Breaking changes

- Malformed CSV files are now rejected instead of being silently misread. An unterminated or misplaced quote previously caused the remainder of the file to be absorbed into a single field, so the import reported success while writing corrupted items. Such files now abort with `Invalid CSV record at line <n>`. Imports that previously appeared to succeed may now fail and require the source file to be corrected. ([#207](https://github.com/Azure/CosmosDBShell/pull/207))

### Fixes

- Local emulator outages are now detected across Cosmos DB commands. Requests fail promptly with an error and return the shell to its disconnected state instead of leaving an unresponsive session labeled as connected.
- A failed or cancelled export no longer destroys its destination file. Exports are written to a temporary file in the destination directory and moved into place only after they complete, so an existing file survives query failures, write failures, and cancellation. An abrupt process termination can leave an unfinished `.cosmos-export-*.tmp` file behind. ([#207](https://github.com/Azure/CosmosDBShell/pull/207))
- `export --max` no longer requests a further query page once the limit is reached, so the reported request charge no longer includes a page whose items were discarded. Query iterators are now disposed. ([#207](https://github.com/Azure/CosmosDBShell/pull/207))
- Shell and MCP command execution is serialized, including nested shell calls, so concurrent requests can no longer interleave and corrupt the shared connection and navigation state. Waiting for a destructive confirmation does not hold the execution lock. ([#207](https://github.com/Azure/CosmosDBShell/pull/207))
- A destructive MCP command is refused when the connection or navigation context changes while its confirmation is pending, including navigating away and back. It previously ran against the changed context. ([#207](https://github.com/Azure/CosmosDBShell/pull/207))
- Echoing an MCP command line no longer fails the command it announces on hosts without an ANSI terminal, which previously reported `Terminal does not support ANSI` instead of running it. ([#207](https://github.com/Azure/CosmosDBShell/pull/207))
- MCP command lines now list positional arguments in the order the command binds them. A destructive confirmation and the recorded history entry previously followed the client's argument order, so `rmdb` could display its `force` flag in place of the database name. A call that supplies a positional argument while omitting an earlier one is now rejected, because the shell cannot express that call and the recorded command would bind differently on replay. ([#207](https://github.com/Azure/CosmosDBShell/pull/207))
- MCP invocations are now saved to the history file as they run and are bounded by the history size limit. They were previously saved only when a later interactive command was entered. On Linux and macOS, the history file is now restricted to its owner, including an existing file that was previously readable by other users. ([#207](https://github.com/Azure/CosmosDBShell/pull/207))

### Build & pipeline

- Added a dependency on CsvHelper 33.1.0 for CSV parsing. ([#207](https://github.com/Azure/CosmosDBShell/pull/207))

## 1.1.209-preview — 2026-08-26

Expand Down
176 changes: 176 additions & 0 deletions CosmosDBShell.Tests/CommandTests/ExportCommandTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,8 @@ namespace CosmosShell.Tests.CommandTests;
using System.Threading;
using System.Threading.Tasks;
using Azure.Data.Cosmos.Shell.Commands;
using Microsoft.Azure.Cosmos;
using NSubstitute;

public class ExportCommandTests
{
Expand Down Expand Up @@ -213,6 +215,180 @@ public async Task WriteCsvAsync_WithNoItems_ProducesEmptyOutput()
Assert.Equal(string.Empty, writer.ToString());
}

[Theory]
[InlineData(0)]
[InlineData(1)]
[InlineData(2)]
public async Task WriteFileAsync_FailurePreservesExistingFileAndRemovesTemporaryFile(int format)
{
var directory = Path.Join(Path.GetTempPath(), Guid.NewGuid().ToString("N"));
Directory.CreateDirectory(directory);
var path = Path.Join(directory, "export.json");
try
{
await File.WriteAllTextAsync(path, "previous export", TestContext.Current.CancellationToken);
await Assert.ThrowsAsync<IOException>(() => ExportCommand.WriteFileAsync(
FailingItemsAsync(), (ExportFormat)format, path, true, CancellationToken.None));
Assert.Equal("previous export", await File.ReadAllTextAsync(path, TestContext.Current.CancellationToken));
Assert.Single(Directory.GetFiles(directory));
}
finally
{
Directory.Delete(directory, recursive: true);
}
}

[Theory]
[InlineData(true)]
[InlineData(false)]
public async Task WriteFileAsync_RespectsOverwriteFlag(bool overwrite)
{
var directory = Path.Join(Path.GetTempPath(), Guid.NewGuid().ToString("N"));
Directory.CreateDirectory(directory);
var path = Path.Join(directory, "export.json");
try
{
await File.WriteAllTextAsync(path, "previous export", TestContext.Current.CancellationToken);
Task<int> ExportAsync() => ExportCommand.WriteFileAsync(
ToAsyncEnumerableAsync(JsonSerializer.SerializeToElement(new { id = "1" })),
ExportFormat.JsonLines, path, overwrite, CancellationToken.None);
if (overwrite)
{
Assert.Equal(1, await ExportAsync());
Assert.Equal("{\"id\":\"1\"}\n", await File.ReadAllTextAsync(path, TestContext.Current.CancellationToken));
}
else
{
await Assert.ThrowsAsync<IOException>(ExportAsync);
Assert.Equal("previous export", await File.ReadAllTextAsync(path, TestContext.Current.CancellationToken));
}

Assert.Single(Directory.GetFiles(directory));
}
finally
{
Directory.Delete(directory, recursive: true);
}
}

[Fact]
public async Task WriteFileAsync_CancellationPreservesExistingFile()
{
var directory = Path.Join(Path.GetTempPath(), Guid.NewGuid().ToString("N"));
Directory.CreateDirectory(directory);
var path = Path.Join(directory, "export.json");
using var cancellation = new CancellationTokenSource();
try
{
await File.WriteAllTextAsync(path, "previous export", TestContext.Current.CancellationToken);
await cancellation.CancelAsync();
await Assert.ThrowsAnyAsync<OperationCanceledException>(() => ExportCommand.WriteFileAsync(
ToAsyncEnumerableAsync(JsonSerializer.SerializeToElement(new { id = "1" })),
ExportFormat.Array, path, true, cancellation.Token));
Assert.Equal("previous export", await File.ReadAllTextAsync(path, TestContext.Current.CancellationToken));
Assert.Single(Directory.GetFiles(directory));
}
finally
{
Directory.Delete(directory, recursive: true);
}
}

[Fact]
public async Task EnumerateAsync_LimitAtPageBoundaryDoesNotFetchAnotherPage()
{
using var iterator = Substitute.For<FeedIterator<JsonElement>>();
var response = Substitute.For<FeedResponse<JsonElement>>();
response.GetEnumerator().Returns(_ => ((IEnumerable<JsonElement>)new[] { JsonSerializer.SerializeToElement(new { id = "1" }) }).GetEnumerator());
response.RequestCharge.Returns(3);
iterator.HasMoreResults.Returns(true);
iterator.ReadNextAsync(Arg.Any<CancellationToken>()).Returns(Task.FromResult(response));
var charge = 0.0;
var count = 0;
await foreach (var item in ExportCommand.EnumerateAsync(iterator, 1, value => charge += value, CancellationToken.None))
{
count++;
}

Assert.Equal(1, count);
Assert.Equal(3, charge);
await iterator.Received(1).ReadNextAsync(Arg.Any<CancellationToken>());
}

[Fact]
public async Task WriteCsvAsync_DoesNotRetainSourceDocuments()
{
using var writer = new StringWriter();
Assert.Equal(200, await ExportCommand.WriteCsvAsync(TransientItemsAsync(), writer, ',', TestContext.Current.CancellationToken));
Assert.Contains("199", writer.ToString());
}

[Fact]
public void DeleteTemporaryFile_DoesNotThrowWhenDirectoryIsMissing()
{
var directory = Directory.CreateTempSubdirectory("cosmos-export-test-");
var path = Path.Join(directory.FullName, "export.tmp");
directory.Delete();

Assert.IsAssignableFrom<IOException>(new DirectoryNotFoundException());
ExportCommand.DeleteTemporaryFile(path);
}

[Fact]
public void DeleteTemporaryFile_DoesNotThrowWhenPathIsDirectory()
{
var directory = Directory.CreateTempSubdirectory("cosmos-export-test-");
try
{
ExportCommand.DeleteTemporaryFile(directory.FullName);
Assert.True(directory.Exists);
}
finally
{
directory.Delete();
}
}

[Fact]
public async Task WriteArrayAsync_FlushesIncrementallyWithoutFlushingEachItem()
{
using var stream = new MemoryStream();
var item = JsonSerializer.SerializeToElement(new { value = new string('x', 1024) });
async IAsyncEnumerable<JsonElement> ItemsAsync()
{
yield return item;
Assert.Equal(0, stream.Length);
for (var index = 0; index < 128; index++)
{
yield return item;
}

Assert.True(stream.Length >= 64 * 1024);
await Task.Yield();
}

Assert.Equal(129, await ExportCommand.WriteArrayAsync(ItemsAsync(), stream, TestContext.Current.CancellationToken));
using var result = JsonDocument.Parse(stream.ToArray());
Assert.Equal(129, result.RootElement.GetArrayLength());
}

private static async IAsyncEnumerable<JsonElement> TransientItemsAsync()
{
for (var index = 0; index < 200; index++)
{
using var document = JsonDocument.Parse($"{{\"id\":{index}}}");
yield return document.RootElement;
await Task.Yield();
}
}

private static async IAsyncEnumerable<JsonElement> FailingItemsAsync()
{
yield return JsonSerializer.SerializeToElement(new { id = "1" });
await Task.Yield();
throw new IOException("simulated read failure");
}

private static async IAsyncEnumerable<JsonElement> ToAsyncEnumerableAsync(params JsonElement[] items)
{
foreach (var item in items)
Expand Down
111 changes: 111 additions & 0 deletions CosmosDBShell.Tests/CommandTests/ImportCommandTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -347,6 +347,117 @@ public void BuildCsvObject_MapsColumnsToStringProperties()
Assert.Equal("Alice", element.GetProperty("name").GetString());
}

[Fact]
public void ReadCsvRecords_ReadsValidRowsBeforeReportingMalformedRecord()
{
using var reader = new StringReader("id,name\n1,Alice\n2,\"unterminated");
using var records = ImportCommand.ReadCsvRecords(reader, ',', TestContext.Current.CancellationToken).GetEnumerator();
Assert.True(records.MoveNext());
Assert.True(records.MoveNext());
Assert.Equal(new[] { "1", "Alice" }, records.Current.Fields);
var error = Assert.Throws<CommandException>(() => records.MoveNext());
Assert.Contains("3", error.Message);
}

[Fact]
public void ReadCsvRecords_CancellationStopsBetweenRecords()
{
using var reader = new StringReader("id,name\n1,Alice\n2,Bob");
using var cancellation = new CancellationTokenSource();
using var records = ImportCommand.ReadCsvRecords(reader, ',', cancellation.Token).GetEnumerator();
Assert.True(records.MoveNext());
cancellation.Cancel();
Assert.ThrowsAny<OperationCanceledException>(() => records.MoveNext());
}

[Fact]
public async Task ReadCsvRecordsAsync_MatchesSynchronousRecordsAndStartLines()
{
const string content = "id,name\n1,\"multi\nline\"\n2,Bob\n";
using var reader = new StringReader(content);
var actual = new List<(int StartLine, List<string> Fields)>();
await foreach (var record in ImportCommand.ReadCsvRecordsAsync(reader, ',', TestContext.Current.CancellationToken))
{
actual.Add(record);
}

var expected = ImportCommand.ParseCsvWithLines(content, ',');
Assert.Equal(expected.Select(r => r.StartLine), actual.Select(r => r.StartLine));
Assert.Equal(expected.Select(r => r.Fields), actual.Select(r => r.Fields));
}

[Fact]
public async Task ReadCsvRecordsAsync_ReportsMalformedRecordWithPhysicalLine()
{
using var reader = new StringReader("id,name\n1,Alice\n2,\"unterminated");
var error = await Assert.ThrowsAsync<CommandException>(async () =>
{
await foreach (var _ in ImportCommand.ReadCsvRecordsAsync(reader, ',', TestContext.Current.CancellationToken))
{
// Enumerate only to drive parsing until the malformed record throws.
}
});
Assert.Contains("3", error.Message);
}

[Fact]
public async Task ReadCsvRecordsAsync_CancellationStopsBetweenRecords()
{
using var reader = new StringReader("id,name\n1,Alice\n2,Bob");
using var cancellation = new CancellationTokenSource();
await using var records = ImportCommand.ReadCsvRecordsAsync(reader, ',', cancellation.Token).GetAsyncEnumerator(TestContext.Current.CancellationToken);
Assert.True(await records.MoveNextAsync());
await cancellation.CancelAsync();
await Assert.ThrowsAnyAsync<OperationCanceledException>(async () => await records.MoveNextAsync());
}

[Fact]
public async Task ReadCsvRecordsAsync_CancellationInterruptsBlockedRecordRead()
{
using var reader = new BlockingAfterPrefixReader("id,name\n1,\"multi");
using var cancellation = new CancellationTokenSource();
await using var records = ImportCommand.ReadCsvRecordsAsync(reader, ',', cancellation.Token).GetAsyncEnumerator(TestContext.Current.CancellationToken);
Assert.True(await records.MoveNextAsync());
var move = records.MoveNextAsync().AsTask();
await reader.Blocked.Task.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken);
Assert.False(move.IsCompleted);

await cancellation.CancelAsync();

await Assert.ThrowsAnyAsync<OperationCanceledException>(() => move.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken));
}

private sealed class BlockingAfterPrefixReader(string prefix) : TextReader
{
private bool prefixReturned;

public TaskCompletionSource Blocked { get; } = new(TaskCreationOptions.RunContinuationsAsynchronously);

public override async ValueTask<int> ReadAsync(Memory<char> buffer, CancellationToken cancellationToken = default)
{
if (!this.prefixReturned)
{
this.prefixReturned = true;
prefix.AsSpan().CopyTo(buffer.Span);
return prefix.Length;
}

this.Blocked.TrySetResult();
await Task.Delay(Timeout.Infinite, cancellationToken);
return 0;
}

public override Task<int> ReadAsync(char[] buffer, int index, int count)
=> this.ReadAsync(buffer.AsMemory(index, count)).AsTask();
}

[Fact]
public void ParseCsvWithLines_SkipsBlankLinesWithoutLosingPhysicalLineNumbers()
{
var records = ImportCommand.ParseCsvWithLines("id,name\n\n\n1,Alice\n", ',');
Assert.Equal(4, records[1].StartLine);
}

[Fact]
public void BuildCsvObject_SingleSegmentPartitionKey_StaysTopLevel()
{
Expand Down
Loading
Loading