Skip to content
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@

### Fixes

- Vector `ORDER BY`, `ORDER BY RANK` relevance ranking, and object-shaped `DISTINCT` projections no longer fail with a continuation-token error. These query pipelines execute successfully but cannot export a resumable token, which was previously reported as a command failure. Such queries now return their documents; through MCP they keep reading until the requested limit instead of stopping after one page, and a truncated result is reported as `resultIncomplete` rather than as an exhausted result set. ([#219](https://github.com/Azure/CosmosDBShell/issues/219))
- 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))
Expand Down
285 changes: 285 additions & 0 deletions CosmosDBShell.Tests/CommandTests/QueryCommandTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -5,11 +5,16 @@
namespace CosmosShell.Tests.CommandTests;

using System.Globalization;
using System.Net;
using System.Text;
using System.Text.Json;
using Azure.Data.Cosmos.Shell.Commands;
using Azure.Data.Cosmos.Shell.Core;
using Microsoft.Azure.Cosmos;
using NSubstitute;
using Spectre.Console;

[Collection(CosmosShell.Tests.Shell.ThemeStateTestCollection.Name)]
public class QueryCommandTests
{
private class TestServerSideMetrics : ServerSideMetrics
Expand Down Expand Up @@ -544,4 +549,284 @@ public void ParseIndexPlan_RecognizedEmptyGroups_ReturnsAvailable()
Assert.Empty(utilized);
Assert.Empty(potential);
}

[Fact]
public void TryReadContinuationToken_ResponseWithoutToken_ReportsSupported()
{
using var response = new ResponseMessage(HttpStatusCode.OK);

Assert.True(QueryCommand.TryReadContinuationToken(response, out var continuationToken));
Assert.Null(continuationToken);
}

[Fact]
public void TryReadContinuationToken_ResponseWithToken_ReturnsToken()
{
using var response = new PageResponse("{\"_count\":0,\"Documents\":[]}", () => "next-page");

Assert.True(QueryCommand.TryReadContinuationToken(response, out var continuationToken));
Assert.Equal("next-page", continuationToken);
}

[Theory]
[InlineData("Continuation tokens are not supported for the non streaming order by pipeline.")]
[InlineData("Continuation tokens are not supported by hybrid search.")]
[InlineData("DISTINCT queries only return continuation tokens when there is a matching ORDER BY clause.")]
public void TryReadContinuationToken_PipelineWithoutTokenSupport_ReportsUnsupported(string message)
{
using var response = new PageResponse("{\"_count\":0,\"Documents\":[]}", () => throw new ArgumentException(message));

Assert.False(QueryCommand.TryReadContinuationToken(response, out var continuationToken));
Assert.Null(continuationToken);
}

[Fact]
public async Task ExecuteQueryAsync_PipelineWithoutTokenSupport_ReturnsDocuments()
{
using var shell = ShellInterpreter.CreateInstance();
using var iterator = new FakeFeedIterator(
NonResumablePage("Continuation tokens are not supported for the non streaming order by pipeline.", 1.5, "1", "2", "6"));
var container = CreateContainer(iterator);
var command = new QueryCommand { Query = "SELECT TOP 3 c.id FROM c ORDER BY VectorDistance(c.embedding, [1,0,0])", Max = 10 };

var result = await command.ExecuteQueryAsync(container, shell, CancellationToken.None);

Assert.Equal(["1", "2", "6"], ReadIds(result));
Assert.Null(result.ContinuationToken);
Assert.False(result.IncompleteWithoutContinuation);
Assert.Equal(1.5, result.RequestCharge);
}

[Fact]
public async Task ExecuteQueryAsync_PipelineWithoutTokenSupport_AdvancesThroughEmptyPages()
{
using var shell = ShellInterpreter.CreateInstance();
using var iterator = new FakeFeedIterator(
NonResumablePage("Continuation tokens are not supported by hybrid search.", 1, "1"),
NonResumablePage("Continuation tokens are not supported by hybrid search.", 2),
NonResumablePage("Continuation tokens are not supported by hybrid search.", 3, "2", "6"));
var container = CreateContainer(iterator);
var command = new QueryCommand { Query = "SELECT TOP 3 c.id FROM c ORDER BY RANK FullTextScore(c.text, \"cosmos\")", Max = 10, IsMcpRequest = true };

var result = await command.ExecuteQueryAsync(container, shell, CancellationToken.None);

Assert.Equal(["1", "2", "6"], ReadIds(result));
Assert.Equal(3, iterator.ReadCount);
Assert.Null(result.ContinuationToken);
Assert.False(result.IncompleteWithoutContinuation);
Assert.Equal(6, result.RequestCharge);
}

[Fact]
public async Task ExecuteQueryAsync_PipelineWithoutTokenSupport_ReportsIncompleteResultAtLimit()
{
using var shell = ShellInterpreter.CreateInstance();
using var iterator = new FakeFeedIterator(
NonResumablePage("Continuation tokens are not supported by hybrid search.", 1, "1", "2"),
NonResumablePage("Continuation tokens are not supported by hybrid search.", 1, "6"));
var container = CreateContainer(iterator);
var command = new QueryCommand { Query = "SELECT c.id FROM c ORDER BY RANK FullTextScore(c.text, \"cosmos\")", Max = 2, IsMcpRequest = true };

var output = await CaptureConsoleAsync(() => command.ExecuteQueryAsync(container, shell, CancellationToken.None));

Assert.Equal(["1", "2"], ReadIds(output.Result));
Assert.Null(output.Result.ContinuationToken);
Assert.True(output.Result.IncompleteWithoutContinuation);
Assert.Contains("cannot be resumed", output.Text);
}

[Fact]
public async Task ExecuteQueryAsync_PipelineWithoutTokenSupport_CancelledBetweenPages_ReportsIncompleteResult()
{
using var shell = ShellInterpreter.CreateInstance();
using var cancellation = new CancellationTokenSource();
using var iterator = new FakeFeedIterator(
NonResumablePage("Continuation tokens are not supported by hybrid search.", 1, "1"),
NonResumablePage("Continuation tokens are not supported by hybrid search.", 1, "2"));
iterator.AfterRead = cancellation.Cancel;
var container = CreateContainer(iterator);
var command = new QueryCommand { Query = "SELECT c.id FROM c ORDER BY RANK FullTextScore(c.text, \"cosmos\")", Max = 10, IsMcpRequest = true };

var result = await command.ExecuteQueryAsync(container, shell, cancellation.Token);

Assert.Equal(["1"], ReadIds(result));
Assert.Equal(1, iterator.ReadCount);
Assert.Null(result.ContinuationToken);
Assert.True(result.IncompleteWithoutContinuation);
}

[Fact]
public async Task ExecuteQueryAsync_ObjectShapedDistinctWithOrderBy_ReturnsDocumentsWithoutToken()
{
using var shell = ShellInterpreter.CreateInstance();
using var iterator = new FakeFeedIterator(
NonResumableDocumentPage(
"DISTINCT queries only return continuation tokens when there is a matching ORDER BY clause.",
1,
"{\"category\":\"A\"}",
"{\"category\":\"B\"}",
"{\"category\":\"C\"}"));
var container = CreateContainer(iterator);
var command = new QueryCommand { Query = "SELECT DISTINCT c.category FROM c ORDER BY c.category", Max = 10, IsMcpRequest = true };

var result = await command.ExecuteQueryAsync(container, shell, CancellationToken.None);

Assert.Equal(["A", "B", "C"], ReadValues(result, "category"));
Assert.Null(result.ContinuationToken);
Assert.False(result.IncompleteWithoutContinuation);
}

[Fact]
public async Task ExecuteQueryAsync_TokenExportRefusedOnLaterPage_DiscardsEarlierToken()
{
using var shell = ShellInterpreter.CreateInstance();
using var iterator = new FakeFeedIterator(
ResumablePage("stale-token", 1, "1"),
NonResumablePage("Continuation tokens are not supported by hybrid search.", 1, "2"));
var container = CreateContainer(iterator);
var command = new QueryCommand { Query = "SELECT c.id FROM c", Max = 10 };

var result = await command.ExecuteQueryAsync(container, shell, CancellationToken.None);

Assert.Equal(["1", "2"], ReadIds(result));
Assert.Equal(2, iterator.ReadCount);
Assert.Null(result.ContinuationToken);
Assert.False(result.IncompleteWithoutContinuation);
}

[Fact]
public async Task ExecuteQueryAsync_ResumablePage_KeepsSingleMcpPageAndToken()
{
using var shell = ShellInterpreter.CreateInstance();
using var iterator = new FakeFeedIterator(
ResumablePage("next-page", 1, "A", "B"),
ResumablePage(null, 1, "C"));
var container = CreateContainer(iterator);
var command = new QueryCommand { Query = "SELECT DISTINCT VALUE c.category FROM c ORDER BY c.category", Max = 10, IsMcpRequest = true };
Comment thread
mkrueger marked this conversation as resolved.

var result = await command.ExecuteQueryAsync(container, shell, CancellationToken.None);

Assert.Equal(1, iterator.ReadCount);
Assert.Equal("next-page", result.ContinuationToken);
Assert.False(result.IncompleteWithoutContinuation);
}

[Fact]
public async Task ExecuteQueryAsync_ResumableQueryAtLimit_ReportsLimitWithoutResumeWarning()
{
using var shell = ShellInterpreter.CreateInstance();
using var iterator = new FakeFeedIterator(
ResumablePage("next-page", 1, "1", "2", "3"));
var container = CreateContainer(iterator);
var command = new QueryCommand { Query = "SELECT * FROM c", Max = 2 };

var output = await CaptureConsoleAsync(() => command.ExecuteQueryAsync(container, shell, CancellationToken.None));

Assert.Equal(["1", "2"], ReadIds(output.Result));
Assert.Equal("next-page", output.Result.ContinuationToken);
Assert.False(output.Result.IncompleteWithoutContinuation);
Assert.Contains("Results limited to 2 items", output.Text);
Assert.DoesNotContain("cannot be resumed", output.Text);
}

private static Container CreateContainer(FeedIterator iterator)
{
var container = Substitute.For<Container>();
container.GetItemQueryStreamIterator(Arg.Any<string>(), Arg.Any<string>(), Arg.Any<QueryRequestOptions>()).Returns(iterator);
return container;
}

private static ResponseMessage NonResumablePage(string message, double requestCharge, params string[] ids)
{
return NonResumableDocumentPage(message, requestCharge, [.. ids.Select(id => $"{{\"id\":\"{id}\"}}")]);
}

private static ResponseMessage NonResumableDocumentPage(string message, double requestCharge, params string[] documents)
{
return CreatePage(requestCharge, () => throw new ArgumentException(message), documents);
}

private static ResponseMessage ResumablePage(string? continuationToken, double requestCharge, params string[] ids)
{
return CreatePage(requestCharge, () => continuationToken, [.. ids.Select(id => $"{{\"id\":\"{id}\"}}")]);
}

private static ResponseMessage CreatePage(double requestCharge, Func<string?> continuationToken, string[] documents)
{
var response = new PageResponse($"{{\"_count\":{documents.Length},\"Documents\":[{string.Join(",", documents)}]}}", continuationToken);
response.Headers.Add("x-ms-request-charge", requestCharge.ToString(CultureInfo.InvariantCulture));
return response;
}

private static string[] ReadIds(CommandState state)
{
return ReadValues(state, "id");
}

private static string[] ReadValues(CommandState state, string property)
{
using var document = JsonDocument.Parse(state.GenerateOutputText());
return [.. document.RootElement.GetProperty("values").EnumerateArray().Select(value => value.GetProperty(property).GetString()!)];
}

private static async Task<(CommandState Result, string Text)> CaptureConsoleAsync(Func<Task<CommandState>> action)
{
var saved = AnsiConsole.Console;
using var writer = new StringWriter();
try
{
AnsiConsole.Console = AnsiConsole.Create(new AnsiConsoleSettings
{
Ansi = AnsiSupport.No,
ColorSystem = ColorSystemSupport.NoColors,
Out = new AnsiConsoleOutput(writer),
});
AnsiConsole.Console.Profile.Width = 200;

var result = await action();
return (result, writer.ToString());
}
finally
{
AnsiConsole.Console = saved;
}
}

private sealed class PageResponse : ResponseMessage
{
private readonly Func<string?> continuationToken;

public PageResponse(string content, Func<string?> continuationToken)
: base(HttpStatusCode.OK)
{
this.continuationToken = continuationToken;
this.Content = new MemoryStream(Encoding.UTF8.GetBytes(content));
}

public override string ContinuationToken => this.continuationToken()!;
}

private sealed class FakeFeedIterator : FeedIterator
{
private readonly Queue<ResponseMessage> pages;

public FakeFeedIterator(params ResponseMessage[] pages)
{
this.pages = new Queue<ResponseMessage>(pages);
}

public int ReadCount { get; private set; }

public Action? AfterRead { get; set; }

public override bool HasMoreResults => this.pages.Count > 0;

public override Task<ResponseMessage> ReadNextAsync(CancellationToken cancellationToken = default)
{
this.ReadCount++;
var page = this.pages.Dequeue();
this.AfterRead?.Invoke();
return Task.FromResult(page);
}
}
}
33 changes: 33 additions & 0 deletions CosmosDBShell.Tests/McpResponseFactoryTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -249,4 +249,37 @@ public void CreateSuccess_OmitsRequestChargeWhenNotSet()
Assert.NotNull(result.StructuredContent);
Assert.False(result.StructuredContent!.Value.TryGetProperty("requestCharge", out _));
}

[Fact]
public void CreateSuccess_IncompletePage_MarksResultIncomplete()
{
var commandState = new CommandState
{
IsPage = true,
IncompleteWithoutContinuation = true,
Result = new ShellJson(JsonSerializer.SerializeToElement(new { result = "success" })),
};

var result = McpResponseFactory.CreateSuccess(commandState, new ConnectedState(null!));

Assert.NotNull(result.StructuredContent);
var structured = result.StructuredContent!.Value;
Assert.Equal(JsonValueKind.Null, structured.GetProperty("continuationToken").ValueKind);
Assert.True(structured.GetProperty("resultIncomplete").GetBoolean());
}

[Fact]
public void CreateSuccess_CompletePage_OmitsResultIncomplete()
{
var commandState = new CommandState
{
IsPage = true,
Result = new ShellJson(JsonSerializer.SerializeToElement(new { result = "success" })),
};

var result = McpResponseFactory.CreateSuccess(commandState, new ConnectedState(null!));

Assert.NotNull(result.StructuredContent);
Assert.False(result.StructuredContent!.Value.TryGetProperty("resultIncomplete", out _));
}
}
1 change: 1 addition & 0 deletions CosmosDBShell.Tests/ToolOperationsTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ public void GetTool_PagedMaxDescription_DocumentsSinglePageSemanticsWithoutChang
var maxDescription = schema.GetProperty("properties").GetProperty("max").GetProperty("description").GetString();

Assert.Contains("continuationToken", maxDescription);
Assert.Contains("resultIncomplete", maxDescription);

var shellDescription = factory.Options.Single(option => option.Name[0] == "max").GetDescription(commandName);
Assert.DoesNotContain("continuationToken", shellDescription);
Expand Down
Loading
Loading