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
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,11 @@ All notable changes to this project will be documented in this file.
The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/),
and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).

## [4.0.18] - 2026-05-15

### Fixed
- Queries without `ORDER BY` now return documents in insertion order, matching real Cosmos DB behavior (Issue #72). Previously documents were returned in hash-map order due to `ConcurrentDictionary` enumeration. Added insertion-order tracking to `InMemoryContainer` that maintains document position across create, replace, upsert, and delete operations.

## [4.0.17] - 2026-05-14

### Fixed
Expand Down
102 changes: 94 additions & 8 deletions src/CosmosDB.InMemoryEmulator/InMemoryContainer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,11 @@ internal class InMemoryContainer : Container, IContainerTestSetup
private readonly ConcurrentDictionary<(string Id, string PartitionKey), string> _items = new();
private readonly ConcurrentDictionary<(string Id, string PartitionKey), string> _etags = new();
private readonly ConcurrentDictionary<(string Id, string PartitionKey), DateTimeOffset> _timestamps = new();
// Ref: Observed behavior on Windows Cosmos DB Emulator — documents return in
// insertion order when no ORDER BY is applied. Maintains creation-time ordering
// so GetAllItemsForPartition can enumerate deterministically.
private readonly List<(string Id, string PartitionKey)> _insertionOrder = new();
private readonly object _insertionOrderLock = new();
private readonly List<(DateTimeOffset Timestamp, string Id, string PartitionKey, string Json, bool IsDelete)> _changeFeed = new();
private readonly object _changeFeedLock = new();
private long _changeFeedLsnCounter;
Expand Down Expand Up @@ -357,6 +362,7 @@ public void ClearItems()
_items.Clear();
_etags.Clear();
_timestamps.Clear();
lock (_insertionOrderLock) { _insertionOrder.Clear(); }
lock (_changeFeedLock) { _changeFeed.Clear(); }
}

Expand Down Expand Up @@ -402,9 +408,11 @@ public void LoadPersistedState()
/// </summary>
public string ExportState()
{
var items = _items
.Where(kvp => !IsExpired(kvp.Key))
.Select(kvp => JsonParseHelpers.ParseJson(kvp.Value)).ToList();
List<(string Id, string PartitionKey)> orderedKeys;
lock (_insertionOrderLock) { orderedKeys = _insertionOrder.ToList(); }
var items = orderedKeys
.Where(key => _items.ContainsKey(key) && !IsExpired(key))
.Select(key => JsonParseHelpers.ParseJson(_items[key])).ToList();
var state = new JObject { ["items"] = new JArray(items) };
return state.ToString(Formatting.Indented);
}
Expand Down Expand Up @@ -437,6 +445,7 @@ public void ImportState(string json)
_etags[key] = importEtag;
_timestamps[key] = DateTimeOffset.UtcNow;
_items[key] = EnrichWithSystemProperties(itemJson, importEtag, _timestamps[key]);
lock (_insertionOrderLock) { _insertionOrder.Add(key); }
}
}

Expand Down Expand Up @@ -489,6 +498,25 @@ public void RestoreToPointInTime(DateTimeOffset pointInTime)
_etags.Clear();
_timestamps.Clear();
_itemLocks.Clear();
lock (_insertionOrderLock) { _insertionOrder.Clear(); }

// Rebuild insertion order from the change feed replay sequence.
// Track which keys were first created (not deleted) to preserve original insertion order.
var insertionOrderKeys = new List<(string Id, string PartitionKey)>();
foreach (var entry in feedSnapshot)
{
if (entry.IsDelete && entry.Json.Contains("\"_ttlEviction\":true", StringComparison.Ordinal))
continue;
var entryKey = (entry.Id, entry.PartitionKey);
if (entry.IsDelete)
{
insertionOrderKeys.Remove(entryKey);
}
else if (!insertionOrderKeys.Contains(entryKey))
{
insertionOrderKeys.Add(entryKey);
}
}

foreach (var kvp in lastPerKey)
{
Expand All @@ -500,6 +528,15 @@ public void RestoreToPointInTime(DateTimeOffset pointInTime)
_etags[key] = etag;
_timestamps[key] = pointInTime;
}

lock (_insertionOrderLock)
{
foreach (var key in insertionOrderKeys)
{
if (_items.ContainsKey(key))
_insertionOrder.Add(key);
}
}
}

// ─── IndexingPolicy ───────────────────────────────────────────────────────
Expand Down Expand Up @@ -738,6 +775,7 @@ public override async Task<ItemResponse<T>> CreateItemAsync<T>(
}

TrackBatchWrite(key);
lock (_insertionOrderLock) { _insertionOrder.Add(key); }
var etag = GenerateETag();
_etags[key] = etag;
_timestamps[key] = DateTimeOffset.UtcNow;
Expand Down Expand Up @@ -773,6 +811,7 @@ public override async Task<ItemResponse<T>> CreateItemAsync<T>(
_items.TryRemove(key, out _);
_etags.TryRemove(key, out _);
_timestamps.TryRemove(key, out _);
lock (_insertionOrderLock) { _insertionOrder.Remove(key); }
throw;
}
}
Expand Down Expand Up @@ -867,6 +906,7 @@ public override async Task<ItemResponse<T>> UpsertItemAsync<T>(
}

TrackBatchWrite(key);
if (!existed) { lock (_insertionOrderLock) { _insertionOrder.Add(key); } }

try
{
Expand All @@ -893,6 +933,7 @@ public override async Task<ItemResponse<T>> UpsertItemAsync<T>(
_items.TryRemove(key, out _);
_etags.TryRemove(key, out _);
_timestamps.TryRemove(key, out _);
lock (_insertionOrderLock) { _insertionOrder.Remove(key); }
}
throw;
}
Expand Down Expand Up @@ -1045,6 +1086,7 @@ public override async Task<ItemResponse<T>> DeleteItemAsync<T>(
_items.TryRemove(key, out _);
_etags.TryRemove(key, out _);
_timestamps.TryRemove(key, out _);
lock (_insertionOrderLock) { _insertionOrder.Remove(key); }

TrackBatchWrite(key);
try
Expand All @@ -1057,6 +1099,7 @@ public override async Task<ItemResponse<T>> DeleteItemAsync<T>(
_items[key] = existingJson;
if (previousEtag is not null) _etags[key] = previousEtag;
if (previousTimestamp.HasValue) _timestamps[key] = previousTimestamp.Value;
lock (_insertionOrderLock) { _insertionOrder.Add(key); }
throw;
}

Expand Down Expand Up @@ -1236,6 +1279,7 @@ public override async Task<ResponseMessage> CreateItemStreamAsync(
}

TrackBatchWrite(key);
lock (_insertionOrderLock) { _insertionOrder.Add(key); }
var etag = GenerateETag();
_etags[key] = etag;
_timestamps[key] = DateTimeOffset.UtcNow;
Expand All @@ -1252,6 +1296,7 @@ public override async Task<ResponseMessage> CreateItemStreamAsync(
_items.TryRemove(key, out _);
_etags.TryRemove(key, out _);
_timestamps.TryRemove(key, out _);
lock (_insertionOrderLock) { _insertionOrder.Remove(key); }
throw;
}

Expand Down Expand Up @@ -1360,6 +1405,7 @@ public override async Task<ResponseMessage> UpsertItemStreamAsync(
}

TrackBatchWrite(key);
if (!existed) { lock (_insertionOrderLock) { _insertionOrder.Add(key); } }
try
{
ExecutePostTriggers(requestOptions, JsonParseHelpers.ParseJson(enrichedJson), "Upsert");
Expand All @@ -1378,6 +1424,7 @@ public override async Task<ResponseMessage> UpsertItemStreamAsync(
_items.TryRemove(key, out _);
_etags.TryRemove(key, out _);
_timestamps.TryRemove(key, out _);
lock (_insertionOrderLock) { _insertionOrder.Remove(key); }
}
throw;
}
Expand Down Expand Up @@ -1525,6 +1572,7 @@ public override async Task<ResponseMessage> DeleteItemStreamAsync(
_items.TryRemove(key, out _);
_etags.TryRemove(key, out _);
_timestamps.TryRemove(key, out _);
lock (_insertionOrderLock) { _insertionOrder.Remove(key); }

TrackBatchWrite(key);
try
Expand All @@ -1537,6 +1585,7 @@ public override async Task<ResponseMessage> DeleteItemStreamAsync(
_items[key] = existingJson;
if (previousEtag is not null) _etags[key] = previousEtag;
if (previousTimestamp.HasValue) _timestamps[key] = previousTimestamp.Value;
lock (_insertionOrderLock) { _insertionOrder.Add(key); }
throw;
}

Expand Down Expand Up @@ -2364,6 +2413,7 @@ public override Task<ContainerResponse> DeleteContainerAsync(
_items.Clear();
_etags.Clear();
_timestamps.Clear();
lock (_insertionOrderLock) { _insertionOrder.Clear(); }
lock (_changeFeedLock) { _changeFeed.Clear(); }
_storedProcedures.Clear();
_userDefinedFunctions.Clear();
Expand All @@ -2386,6 +2436,7 @@ public override Task<ResponseMessage> DeleteContainerStreamAsync(
_items.Clear();
_etags.Clear();
_timestamps.Clear();
lock (_insertionOrderLock) { _insertionOrder.Clear(); }
lock (_changeFeedLock) { _changeFeed.Clear(); }
_storedProcedures.Clear();
_userDefinedFunctions.Clear();
Expand Down Expand Up @@ -2449,6 +2500,7 @@ public override Task<ResponseMessage> DeleteAllItemsByPartitionKeyStreamAsync(
_items.TryRemove(key, out _);
_etags.TryRemove(key, out _);
_timestamps.TryRemove(key, out _);
lock (_insertionOrderLock) { _insertionOrder.Remove(key); }
RecordDeleteTombstone(key.Id, pk, partitionKey);
}
return Task.FromResult(CreateResponseMessage(HttpStatusCode.OK));
Expand Down Expand Up @@ -3408,6 +3460,7 @@ private void EvictIfExpired((string Id, string PartitionKey) key)
_items.TryRemove(key, out _);
_etags.TryRemove(key, out _);
_timestamps.TryRemove(key, out _);
lock (_insertionOrderLock) { _insertionOrder.Remove(key); }

// Record a delete tombstone in the change feed so consumers see TTL evictions
RecordDeleteTombstone(key.Id, key.PartitionKey, isTtlEviction: true);
Expand Down Expand Up @@ -3467,6 +3520,27 @@ internal void RestoreSnapshot(
_items.TryRemove(key, out _);
_etags.TryRemove(key, out _);
_timestamps.TryRemove(key, out _);
lock (_insertionOrderLock) { _insertionOrder.Remove(key); }
}
}

// Restore insertion order for keys that were newly added during the batch
// but didn't exist in the snapshot — they need to be removed from insertion order.
// Keys that existed in the snapshot should already be in insertion order.
lock (_insertionOrderLock)
{
foreach (var key in touchedKeys)
{
if (!itemsSnapshot.ContainsKey(key))
{
// Already removed above
}
else if (!_insertionOrder.Contains(key))
{
// Item existed in snapshot but is missing from insertion order
// (shouldn't normally happen, but be safe)
_insertionOrder.Add(key);
}
}
}

Expand Down Expand Up @@ -4061,8 +4135,16 @@ private static void EnrichStoredProcedureSystemProperties(StoredProcedurePropert
// Private helpers — Query execution pipeline
// ═══════════════════════════════════════════════════════════════════════════

// Ref: Observed behavior on the Windows Cosmos DB Emulator (priority 6):
// Documents returned by queries without ORDER BY are in insertion order.
private IEnumerable<string> GetAllItemsForPartition(QueryRequestOptions requestOptions)
{
List<(string Id, string PartitionKey)> orderedKeys;
lock (_insertionOrderLock)
{
orderedKeys = _insertionOrder.ToList();
}

if (requestOptions?.PartitionKey is not null
&& requestOptions.PartitionKey != PartitionKey.None)
{
Expand All @@ -4076,15 +4158,19 @@ private IEnumerable<string> GetAllItemsForPartition(QueryRequestOptions requestO
if (queryComponents > 0 && queryComponents < PartitionKeyPaths.Count)
{
var prefix = pk + "|";
return _items
.Where(kvp => (kvp.Key.PartitionKey?.StartsWith(prefix, StringComparison.Ordinal) ?? false) && !IsExpired(kvp.Key))
.Select(kvp => kvp.Value);
return orderedKeys
.Where(key => (key.PartitionKey?.StartsWith(prefix, StringComparison.Ordinal) ?? false) && !IsExpired(key) && _items.ContainsKey(key))
.Select(key => _items[key]);
}
}

return _items.Where(kvp => kvp.Key.PartitionKey == pk && !IsExpired(kvp.Key)).Select(kvp => kvp.Value);
return orderedKeys
.Where(key => key.PartitionKey == pk && !IsExpired(key) && _items.ContainsKey(key))
.Select(key => _items[key]);
}
return _items.Where(kvp => !IsExpired(kvp.Key)).Select(kvp => kvp.Value);
return orderedKeys
.Where(key => !IsExpired(key) && _items.ContainsKey(key))
.Select(key => _items[key]);
}

private static int CountPartitionKeyComponents(PartitionKey partitionKey)
Expand Down
2 changes: 1 addition & 1 deletion src/Directory.Build.props
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@
<!-- Shared NuGet package metadata for all source (packable) projects -->
<PropertyGroup>
<IsPackable>true</IsPackable>
<Version>4.0.17</Version>
<Version>4.0.18</Version>
<Authors>lemonlion</Authors>
<Copyright>Copyright (c) 2026 lemonlion</Copyright>
<PackageLicenseExpression>MIT</PackageLicenseExpression>
Expand Down
Loading
Loading