diff --git a/README.md b/README.md
index b572b2c..637df1a 100644
--- a/README.md
+++ b/README.md
@@ -225,8 +225,10 @@ the others. Running queries finish before their slots are reassigned.
SQL Server work splits into databases; Oracle work splits into schemas and then
individual objects, including packages, package bodies, and database links.
One Oracle server with one schema can therefore use all four slots for object
-DDL. Oracle sessions are reused exclusively by one job at a time, and schema
-metadata is loaded once. Set `--max-parallelism 1` for sequential extraction
+DDL. Each Oracle job returns its connection to the provider pool before releasing
+its scheduler slot, keeping checked-out sessions within the shared budget even
+across many server entries. Schema metadata is loaded once. Set
+`--max-parallelism 1` for sequential extraction
across all engines. Linked-server discovery still runs in depth rounds using
the same shared budget.
Use `--skip-metrics` when only object definitions are needed; it skips the
diff --git a/cli/src/SyncSql.Extraction.Oracle/OracleExtractionConnections.cs b/cli/src/SyncSql.Extraction.Oracle/OracleExtractionConnections.cs
index 55d1112..3aa9a7a 100644
--- a/cli/src/SyncSql.Extraction.Oracle/OracleExtractionConnections.cs
+++ b/cli/src/SyncSql.Extraction.Oracle/OracleExtractionConnections.cs
@@ -1,36 +1,19 @@
-using System.Collections.Concurrent;
-using System.Data.Common;
+using System.Data.Common;
namespace SyncSql.Extraction.Oracle;
-/// Exclusive sessions reused by scheduled jobs, with metadata transforms initialized once.
-internal sealed class OracleExtractionConnections(Func create) : IAsyncDisposable
+/// Each scheduled job owns a connection lease and releases it before returning its scheduler slot.
+internal sealed class OracleExtractionConnections(Func create)
{
- private readonly ConcurrentBag _idle = [];
-
public async Task UseAsync(Func> work, CancellationToken cancellationToken)
{
cancellationToken.ThrowIfCancellationRequested();
- if (!_idle.TryTake(out DbConnection? connection))
- {
- connection = create();
- try
- {
- await connection.OpenAsync(cancellationToken);
- await OracleConnectionFactory.InitializeAsync(connection, cancellationToken);
- }
- catch
- {
- await connection.DisposeAsync();
- throw;
- }
- }
- try { return await work(connection); }
- finally { _idle.Add(connection); }
- }
-
- public async ValueTask DisposeAsync()
- {
- while (_idle.TryTake(out DbConnection? connection)) { await connection.DisposeAsync(); }
+ // ODP.NET owns physical-session pooling. Keeping open connections here would retain
+ // pool leases while server/schema parents wait for queued children, potentially leaving
+ // every scheduler slot blocked in OpenAsync with no slot available to drain those children.
+ await using DbConnection connection = create();
+ await connection.OpenAsync(cancellationToken);
+ await OracleConnectionFactory.InitializeAsync(connection, cancellationToken);
+ return await work(connection);
}
}
diff --git a/cli/src/SyncSql.Extraction.Oracle/OracleObjectExtractor.cs b/cli/src/SyncSql.Extraction.Oracle/OracleObjectExtractor.cs
index a5239b4..c6de315 100644
--- a/cli/src/SyncSql.Extraction.Oracle/OracleObjectExtractor.cs
+++ b/cli/src/SyncSql.Extraction.Oracle/OracleObjectExtractor.cs
@@ -46,7 +46,7 @@ public async Task ExtractAsync(ServerConfig server, Effective
string serviceName = server.ServiceName
?? throw new InvalidOperationException($"Oracle server '{server.Name}' is missing required key 'serviceName'.");
ExtractionProgressAggregator progress = new(options.Progress);
- await using OracleExtractionConnections connections = new(() => _createConnection(server, options.Credentials));
+ OracleExtractionConnections connections = new(() => _createConnection(server, options.Credentials));
progress.Report(new($"Connecting to {serviceName}"));
var (owners, links) = await connections.UseAsync(async connection =>
{
diff --git a/cli/tests/SyncSql.Extraction.Oracle.Tests/FakeOracleDatabase.cs b/cli/tests/SyncSql.Extraction.Oracle.Tests/FakeOracleDatabase.cs
index 51de81a..015ef65 100644
--- a/cli/tests/SyncSql.Extraction.Oracle.Tests/FakeOracleDatabase.cs
+++ b/cli/tests/SyncSql.Extraction.Oracle.Tests/FakeOracleDatabase.cs
@@ -15,6 +15,8 @@ internal sealed class FakeOracleDatabase : DbConnection
public List Queries { get; } = [];
public bool WasDisposed { get; private set; }
public Exception? OpenFailure { get; init; }
+ public Func? BeforeOpenAsync { get; init; }
+ public Action? OnDispose { get; init; }
[AllowNull] public override string ConnectionString { get; set; } = "";
public override string Database => "APP";
public override string DataSource => "fake";
@@ -25,11 +27,22 @@ public override void Open()
if (OpenFailure is { } failure) { throw failure; }
_state = ConnectionState.Open;
}
+ public override async Task OpenAsync(CancellationToken cancellationToken)
+ {
+ cancellationToken.ThrowIfCancellationRequested();
+ if (BeforeOpenAsync is { } beforeOpen) { await beforeOpen(cancellationToken); }
+ Open();
+ }
public override void Close() => _state = ConnectionState.Closed;
public override void ChangeDatabase(string databaseName) => throw new NotSupportedException();
protected override DbTransaction BeginDbTransaction(IsolationLevel isolationLevel) => throw new NotSupportedException();
protected override DbCommand CreateDbCommand() => new FakeCommand(this);
- protected override void Dispose(bool disposing) { WasDisposed = true; base.Dispose(disposing); }
+ protected override void Dispose(bool disposing)
+ {
+ if (!WasDisposed) { OnDispose?.Invoke(); }
+ WasDisposed = true;
+ base.Dispose(disposing);
+ }
public static OracleException Error(int number) => (OracleException)typeof(OracleException)
.GetConstructor(BindingFlags.Instance | BindingFlags.NonPublic, null, [typeof(int), typeof(string), typeof(string), typeof(string), typeof(int)], null)!
diff --git a/cli/tests/SyncSql.Extraction.Oracle.Tests/OracleExtractionConnectionsTests.cs b/cli/tests/SyncSql.Extraction.Oracle.Tests/OracleExtractionConnectionsTests.cs
index ad1080d..bcd39a2 100644
--- a/cli/tests/SyncSql.Extraction.Oracle.Tests/OracleExtractionConnectionsTests.cs
+++ b/cli/tests/SyncSql.Extraction.Oracle.Tests/OracleExtractionConnectionsTests.cs
@@ -15,32 +15,50 @@ public async Task FailedSessionIsDisposedAndReplacedBeforeRunningWork(bool failO
OpenFailure = failOpen ? failure : null,
Execute = (_, _) => throw failure,
};
- using FakeOracleDatabase healthy = new() { Execute = (_, _) => null };
+ using FakeOracleDatabase first = new() { Execute = (_, _) => null };
+ using FakeOracleDatabase second = new() { Execute = (_, _) => null };
+ FakeOracleDatabase[] sessions = [failed, first, second];
int created = 0, executed = 0;
- await using (OracleExtractionConnections pool = new(() => ++created == 1 ? failed : healthy))
+ OracleExtractionConnections connections = new(() => sessions[created++]);
+ OracleException actual = await Assert.ThrowsAsync(() => connections.UseAsync(_ =>
{
- OracleException actual = await Assert.ThrowsAsync(() => pool.UseAsync(_ =>
- {
- executed++;
- return Task.FromResult(0);
- }, CancellationToken.None));
- Assert.Same(failure, actual);
- Assert.True(failed.WasDisposed);
- Assert.Equal(0, executed);
+ executed++;
+ return Task.FromResult(0);
+ }, CancellationToken.None));
+ Assert.Same(failure, actual);
+ Assert.True(failed.WasDisposed);
+ Assert.Equal(0, executed);
- for (int i = 0; i < 2; i++)
+ for (int i = 1; i < sessions.Length; i++)
+ {
+ int result = await connections.UseAsync(connection =>
{
- int result = await pool.UseAsync(connection =>
- {
- Assert.Same(healthy, connection);
- return Task.FromResult(++executed);
- }, CancellationToken.None);
- Assert.Equal(i + 1, result);
- }
- Assert.Equal(2, created);
- Assert.Single(healthy.Queries);
- Assert.False(healthy.WasDisposed);
+ Assert.Same(sessions[i], connection);
+ Assert.False(sessions[i].WasDisposed);
+ return Task.FromResult(++executed);
+ }, CancellationToken.None);
+ Assert.Equal(i, result);
+ Assert.Single(sessions[i].Queries);
+ Assert.True(sessions[i].WasDisposed);
}
- Assert.True(healthy.WasDisposed);
+ Assert.Equal(3, created);
+ }
+
+ [Theory]
+ [InlineData(false)]
+ [InlineData(true)]
+ public async Task WorkFailureOrCancellationDisposesLeaseBeforeReturning(bool cancel)
+ {
+ using CancellationTokenSource cancellation = new();
+ using FakeOracleDatabase session = new() { Execute = (_, _) => null };
+ OracleExtractionConnections connections = new(() => session);
+ Exception failure = cancel ? new OperationCanceledException(cancellation.Token) : FakeOracleDatabase.Error(3113);
+ Exception? actual = await Record.ExceptionAsync(() => connections.UseAsync(_ =>
+ {
+ if (cancel) { cancellation.Cancel(); }
+ return Task.FromException(failure);
+ }, cancellation.Token));
+ Assert.Same(failure, actual);
+ Assert.True(session.WasDisposed);
}
}
diff --git a/cli/tests/SyncSql.Extraction.Oracle.Tests/OracleExtractorTests.cs b/cli/tests/SyncSql.Extraction.Oracle.Tests/OracleExtractorTests.cs
index afa67ae..e437471 100644
--- a/cli/tests/SyncSql.Extraction.Oracle.Tests/OracleExtractorTests.cs
+++ b/cli/tests/SyncSql.Extraction.Oracle.Tests/OracleExtractorTests.cs
@@ -202,7 +202,7 @@ public async Task Extract_EmptySelectionAndMissingService()
[InlineData(null, 4, 7)]
[InlineData(1, 1, 1)]
[InlineData(2, 2, 7)]
- public async Task Extract_SchemasAndObjectsShareTheBudgetAndReuseExclusiveSessions(int? limit, int expectedConcurrency, int ownerCount)
+ public async Task Extract_SchemasAndObjectsShareTheBudgetWithExclusiveSessionLeases(int? limit, int expectedConcurrency, int ownerCount)
{
string[] owners = [.. Enumerable.Range(0, ownerCount).Select(i => $"OWNER{i}")];
System.Collections.Concurrent.ConcurrentBag sessions = [];
@@ -258,7 +258,7 @@ public async Task Extract_SchemasAndObjectsShareTheBudgetAndReuseExclusiveSessio
{
await slotsFilled.Task.WaitAsync(TimeSpan.FromSeconds(10));
Assert.Equal(expectedConcurrency, Volatile.Read(ref started));
- Assert.Equal(expectedConcurrency, sessions.Count);
+ Assert.Equal(expectedConcurrency, sessions.Count(session => !session.WasDisposed));
}
finally
{
@@ -283,26 +283,28 @@ public async Task Extract_SchemasAndObjectsShareTheBudgetAndReuseExclusiveSessio
Assert.Equal(ownerCount == 1 ? 0 : 1, sessions.Sum(session => session.Queries.Count(q => q == OracleQueries.AllDatabaseLinks)));
Assert.Equal(ownerCount, sessions.Sum(session => session.Queries.Count(q => q == OracleQueries.AllObjectGrants)));
Assert.Equal(ownerCount, sessions.Sum(session => session.Queries.Count(q => q == OracleQueries.ColumnList)));
- Assert.InRange(sessions.Count, 1, expectedConcurrency);
+ Assert.Equal(1 + ownerCount + expectedObjects, sessions.Count);
Assert.All(sessions, session => Assert.True(session.WasDisposed));
}
[Theory]
[InlineData(false)]
[InlineData(true)]
- public async Task Extract_SequentialWorkReusesDiscoverySession(bool linksOnly)
+ public async Task Extract_SequentialWorkReleasesDiscoveryBeforeOpeningObjectConnections(bool linksOnly)
{
- int created = 0;
- using FakeOracleDatabase db = new() { Execute = Respond };
+ List sessions = [];
OracleObjectExtractor extractor = new(NullLogger.Instance, TimeProvider.System, (_, _) =>
{
- Assert.Equal(1, ++created);
+ Assert.All(sessions, session => Assert.True(session.WasDisposed));
+ FakeOracleDatabase db = new() { Execute = Respond };
+ sessions.Add(db);
return db;
});
ServerConfig server = Server with { ObjectTypes = linksOnly ? ["DatabaseLinks"] : [] };
ExtractionOutcome result = await extractor.ExtractAsync(server, EffectiveFilters.Resolve(null, server),
new ExtractionOptions { Credentials = Credentials, MaxParallelism = 1 }, CancellationToken.None);
Assert.Equal(linksOnly ? 4 : 0, result.Objects.Count);
- Assert.True(db.WasDisposed);
+ Assert.Equal(linksOnly ? 5 : 1, sessions.Count);
+ Assert.All(sessions, session => Assert.True(session.WasDisposed));
}
[Theory]
@@ -360,7 +362,7 @@ public async Task Extract_CancellationOrFailureStopsWorkersAndDisposesConnection
await Assert.ThrowsAnyAsync(() => run.WaitAsync(TimeSpan.FromSeconds(10)));
}
Assert.Equal(2, started);
- Assert.Equal(2, created);
+ Assert.Equal(3, created);
Assert.All(connections, db => Assert.True(db.WasDisposed));
}
finally
@@ -371,4 +373,63 @@ public async Task Extract_CancellationOrFailureStopsWorkersAndDisposesConnection
catch (OracleException) when (fail) { }
}
}
+
+ [Theory]
+ [InlineData(1)]
+ [InlineData(4)]
+ public async Task Extract_QueuedServersAndSchemasDoNotRetainConnectionsOutsideSharedBudget(int limit)
+ {
+ // Model repeated server entries sharing an Oracle pool whose capacity equals the budget.
+ // Retaining any idle discovery/metadata session can strand all slots in OpenAsync.
+ using SemaphoreSlim poolCapacity = new(limit);
+ using CancellationTokenSource cancellation = new(TimeSpan.FromSeconds(10));
+ System.Collections.Concurrent.ConcurrentBag sessions = [];
+ int open = 0, peak = 0;
+ object gate = new();
+ OracleObjectExtractor extractor = new(NullLogger.Instance, TimeProvider.System, (_, _) =>
+ {
+ bool leased = false;
+ FakeOracleDatabase session = new()
+ {
+ BeforeOpenAsync = async token =>
+ {
+ await poolCapacity.WaitAsync(token);
+ leased = true;
+ lock (gate) { peak = Math.Max(peak, ++open); }
+ },
+ OnDispose = () =>
+ {
+ if (!leased) { return; }
+ lock (gate) { open--; }
+ poolCapacity.Release();
+ leased = false;
+ },
+ Execute = (sql, p) => sql == OracleQueries.Schemas
+ ? FakeOracleDatabase.Rows(new { OWNER = "APP" }, new { OWNER = "AUX" }) : Respond(sql, p),
+ };
+ sessions.Add(session);
+ return session;
+ });
+ ServerConfig[] servers = [.. Enumerable.Range(0, 12).Select(i => Server with
+ {
+ Name = $"ORA{i}", ObjectTypes = ["Tables"], ObjectNames = new NameFilter { Exclude = ["^skip$"] },
+ })];
+ ExtractionWorkScheduler scheduler = new(limit);
+ ExtractionOutcome[] results = await scheduler.RunAsync(servers.Select(server =>
+ (DatabaseEngine.Oracle, (Func>)(context =>
+ extractor.ExtractAsync(server, EffectiveFilters.Resolve(null, server),
+ new ExtractionOptions { Credentials = Credentials, WorkContext = context, MaxParallelism = limit },
+ context.CancellationToken)))), cancellation.Token);
+
+ Assert.Equal(servers.Length, results.Length);
+ Assert.All(results, result =>
+ {
+ Assert.Equal(4, result.Objects.Count);
+ Assert.Equal(4, result.MetricsSnapshots.Count);
+ });
+ Assert.InRange(peak, 1, limit);
+ Assert.Equal(0, open);
+ Assert.Equal(limit, poolCapacity.CurrentCount);
+ Assert.All(sessions, session => Assert.True(session.WasDisposed));
+ }
}