This is an automated email from the ASF dual-hosted git repository.
HTHou pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/iotdb-client-csharp.git
The following commit(s) were added to refs/heads/main by this push:
new 178ebb0 Fix SessionPool blocking indefinitely after a server
outage(#62)
178ebb0 is described below
commit 178ebb07bf284033921c606ca3446b64f2cd2eec
Author: CritasWang <[email protected]>
AuthorDate: Fri Jul 31 14:36:51 2026 +0800
Fix SessionPool blocking indefinitely after a server outage(#62)
---
docs/API.md | 2 +-
docs/SessionPool_Exception_Handling.md | 126 ++++++++++++---
src/Apache.IoTDB/ConcurrentClientQueue.cs | 40 ++++-
src/Apache.IoTDB/SessionPool.Builder.cs | 18 ++-
src/Apache.IoTDB/SessionPool.cs | 122 +++++++++++++-
src/Apache.IoTDB/TableSessionPool.Builder.cs | 18 ++-
.../ConcurrentClientQueueTests.cs | 178 +++++++++++++++++++++
.../SessionPoolConfigurationTests.cs | 135 ++++++++++++++++
8 files changed, 598 insertions(+), 41 deletions(-)
diff --git a/docs/API.md b/docs/API.md
index 13bcbc8..43fdd04 100644
--- a/docs/API.md
+++ b/docs/API.md
@@ -44,7 +44,7 @@ var tablet =
| -------------- | ------------------------- | ------------------------ |
----------------------------- |
| Open | bool | open session |
session_pool.Open(false) |
| Close | null | close session |
session_pool.Close() |
-| IsOpen | null | check if session is open |
session_pool.IsOpen() |
+| IsOpen | null | check if the pool was opened
and not yet closed by the caller. It is a lifecycle flag, **not** a
connectivity probe: it stays `true` after the server goes down, because the
client keeps no heartbeat and reconnects on demand instead. |
session_pool.IsOpen() |
| OpenDebugMode | LoggingConfiguration=null | open debug mode |
session_pool.OpenDebugMode() |
| CloseDebugMode | null | close debug mode |
session_pool.CloseDebugMode() |
| SetTimeZone | string | set time zone |
session_pool.GetTimeZone() |
diff --git a/docs/SessionPool_Exception_Handling.md
b/docs/SessionPool_Exception_Handling.md
index 75fd019..99039de 100644
--- a/docs/SessionPool_Exception_Handling.md
+++ b/docs/SessionPool_Exception_Handling.md
@@ -34,9 +34,9 @@ using System;
try
{
var sessionPool = new SessionPool.Builder()
- .Host("127.0.0.1")
- .Port(6667)
- .PoolSize(4)
+ .SetHost("127.0.0.1")
+ .SetPort(6667)
+ .SetPoolSize(4)
.Build();
await sessionPool.Open();
@@ -54,6 +54,52 @@ catch (SessionPoolDepletedException ex)
}
```
+## `IsOpen()` is a lifecycle flag, not a health check
+
+`SessionPool.IsOpen()` reports whether **you** have opened the pool and not
yet closed it. It is not a
+connectivity probe:
+
+- It becomes `true` after a successful `Open()` and only returns to `false`
when you call `Close()`.
+- The client runs no heartbeat, so a server that goes down does **not** flip
it back to `false`.
+ Reconnection happens lazily, on the next operation.
+
+This means the following common guard never re-opens the pool, because the
flag stays `true` forever:
+
+```csharp
+// Anti-pattern: this short-circuits even while every connection is dead
+if (_pool != null && _pool.IsOpen()) return;
+```
+
+To reason about actual availability, use the health metrics below, or simply
let an operation throw
+`SessionPoolDepletedException` and handle it.
+
+## Pool Wait Timeout
+
+Two independent timeouts govern a pool operation:
+
+| Setting | Unit | Default | Controls
|
+| ----------------------------------------- | ---- | ------- |
------------------------------------------------------------------------- |
+| `SetConnectionTimeoutInMs(int)` | ms | 500 | Socket-level
send/receive timeout of an individual connection |
+| `SetPoolWaitTimeoutInMs(int)` | ms | 10000 | How long an
operation waits for a free client before the pool gives up |
+
+```csharp
+var sessionPool = new SessionPool.Builder()
+ .SetHost("127.0.0.1")
+ .SetPort(6667)
+ .SetPoolSize(8)
+ .SetConnectionTimeoutInMs(500) // socket timeout
+ .SetPoolWaitTimeoutInMs(10_000) // give up after 10s of waiting for a
free client
+ .Build();
+```
+
+When the wait budget is exhausted, the operation throws
`SessionPoolDepletedException` with the reason
+`Connection pool is empty and wait time out(...ms)`. Raise
`SetPoolWaitTimeoutInMs` if your workload
+legitimately queues behind long operations; lower it if you would rather fail
fast and retry.
+
+> **Note:** before this setting existed, the wait budget was derived from the
connection timeout and then
+> misinterpreted as seconds, which turned the 500 ms default into a ~41 minute
block. If you are upgrading
+> from an older version and relied on that (unintended) long wait, set
`SetPoolWaitTimeoutInMs` explicitly.
+
## Pool Health Metrics
### Monitoring Pool Status
@@ -62,9 +108,9 @@ The `SessionPool` class exposes real-time health metrics
that can be used for mo
```csharp
var sessionPool = new SessionPool.Builder()
- .Host("127.0.0.1")
- .Port(6667)
- .PoolSize(8)
+ .SetHost("127.0.0.1")
+ .SetPort(6667)
+ .SetPoolSize(8)
.Build();
await sessionPool.Open();
@@ -72,6 +118,7 @@ await sessionPool.Open();
// Check pool health
Console.WriteLine($"Available Clients: {sessionPool.AvailableClients}");
Console.WriteLine($"Total Pool Size: {sessionPool.TotalPoolSize}");
+Console.WriteLine($"Unrealized Capacity: {sessionPool.UnrealizedCapacity}");
Console.WriteLine($"Failed Reconnections: {sessionPool.FailedReconnections}");
```
@@ -81,8 +128,27 @@ Console.WriteLine($"Failed Reconnections:
{sessionPool.FailedReconnections}");
| -------------------- | --------------------- |
------------------------------------------------ | --------------------------- |
| Available Clients | `AvailableClients` | Number of idle clients ready
for use | Alert if < 25% of pool size |
| Total Pool Size | `TotalPoolSize` | Configured maximum pool size
| N/A (constant) |
+| Unrealized Capacity | `UnrealizedCapacity` | Configured capacity currently
holding no connection, refilled on demand | Not an alert signal on its own -
see below |
| Failed Reconnections | `FailedReconnections` | Cumulative count of failed
reconnection attempts | Alert if > 0 and increasing |
+### Capacity is demand-driven
+
+When an operation fails and reconnection also fails, the dead connection is
discarded but its **capacity is
+retained** rather than lost. `UnrealizedCapacity` counts the capacity left
without a connection, and an
+acquisition that finds no idle client materializes one connection before
falling back to waiting.
+Consequences:
+
+- The pool no longer shrinks by one on every failure, so it cannot reach the
state where every caller blocks
+ on a queue nobody will feed.
+- Once the server is reachable again, the pool repopulates itself as load
demands it - no `Close()` +
+ `Open()` cycle is required.
+- **Capacity is refilled on demand, not eagerly.** A connection is only
created when an acquisition finds the
+ idle queue empty. Under sequential or light workloads one connection is
enough to serve every request, so
+ `UnrealizedCapacity` legitimately stays above zero long after the server has
fully recovered. It measures
+ how much of the configured pool has not been materialized, not server
availability.
+- Therefore **do not alert on `UnrealizedCapacity` alone.** Use
`FailedReconnections` to reason about server
+ reachability: it only increases when a reconnection actually fails.
+
## Failure Scenarios and Recovery Strategies
### Scenario 1: Pool Exhaustion (High Load)
@@ -101,9 +167,9 @@ Console.WriteLine($"Failed Reconnections:
{sessionPool.FailedReconnections}");
```csharp
var sessionPool = new SessionPool.Builder()
- .Host("127.0.0.1")
- .Port(6667)
- .PoolSize(16) // Increased from 8
+ .SetHost("127.0.0.1")
+ .SetPort(6667)
+ .SetPoolSize(16) // Increased from 8
.Build();
```
@@ -137,13 +203,25 @@ for (int i = 0; i < maxRetries; i++)
**Symptoms:**
- `SessionPoolDepletedException` with reason "Reconnection failed"
-- `AvailableClients` decreases over time
-- `FailedReconnections` > 0 and increasing
+- `AvailableClients` drops toward 0 while `UnrealizedCapacity` rises
+- `FailedReconnections` > 0 and increasing (this, not `UnrealizedCapacity`, is
the outage signal)
**Root Cause:** IoTDB server unreachable or network issues
**Recovery Strategies:**
+0. **Do nothing but retry.** Capacity is retained and refilled on demand, so
once the server comes back a
+ plain retry succeeds. Reinitialising is only needed if you want to change
configuration or drop
+ accumulated state:
+
+```csharp
+catch (SessionPoolDepletedException ex)
+{
+ // Capacity is retained and refilled on demand - just back off and try
again
+ await Task.Delay(2000);
+}
+```
+
1. **Reinitialize SessionPool:**
```csharp
@@ -159,9 +237,9 @@ catch (SessionPoolDepletedException ex) when
(ex.FailedReconnections > 5)
// Create new pool
sessionPool = new SessionPool.Builder()
- .Host("127.0.0.1")
- .Port(6667)
- .PoolSize(8)
+ .SetHost("127.0.0.1")
+ .SetPort(6667)
+ .SetPoolSize(8)
.Build();
await sessionPool.Open();
@@ -245,9 +323,9 @@ public async Task RateLimitedInsert(string deviceId,
RowRecord record)
```csharp
var sessionPool = new SessionPool.Builder()
- .Host("127.0.0.1")
- .Port(6667)
- .Timeout(120) // Increased timeout for slow server
+ .SetHost("127.0.0.1")
+ .SetPort(6667)
+ .SetConnectionTimeoutInMs(5000) // Increased socket timeout for a slow
server
.Build();
```
@@ -381,10 +459,10 @@ public class ProductionSessionPoolManager
public async Task Initialize()
{
_pool = new SessionPool.Builder()
- .Host("127.0.0.1")
- .Port(6667)
- .PoolSize(8)
- .Timeout(60)
+ .SetHost("127.0.0.1")
+ .SetPort(6667)
+ .SetPoolSize(8)
+ .SetConnectionTimeoutInMs(5000)
.Build();
await _pool.Open();
@@ -481,7 +559,11 @@ public class ProductionSessionPoolManager
The SessionPool exception handling and health monitoring features provide
comprehensive tools for building robust IoTDB applications:
- Use `SessionPoolDepletedException` to understand and react to pool issues
-- Monitor `AvailableClients`, `TotalPoolSize`, and `FailedReconnections`
metrics
+- Treat `IsOpen()` as a lifecycle flag, never as a connectivity check
+- Tune `SetPoolWaitTimeoutInMs` separately from `SetConnectionTimeoutInMs`
+- Monitor `AvailableClients`, `TotalPoolSize`, and `FailedReconnections`; read
`UnrealizedCapacity` as
+ capacity not yet materialized rather than as an outage signal
+- Rely on demand-driven capacity refill for recovery; reinitialise only when
you need to change configuration
- Implement appropriate recovery strategies based on failure scenarios
- Set up proactive monitoring and alerting to prevent issues
- Follow best practices for pool sizing and resource management
diff --git a/src/Apache.IoTDB/ConcurrentClientQueue.cs
b/src/Apache.IoTDB/ConcurrentClientQueue.cs
index 6090c19..f59b631 100644
--- a/src/Apache.IoTDB/ConcurrentClientQueue.cs
+++ b/src/Apache.IoTDB/ConcurrentClientQueue.cs
@@ -58,26 +58,52 @@ namespace Apache.IoTDB
public void AddRef() => Interlocked.Increment(ref _ref);
public int GetRef() => Volatile.Read(ref _ref);
public void RemoveRef() => Interlocked.Decrement(ref _ref);
- public int Timeout { get; set; } = 10;
+
+ /// <summary>
+ /// The maximum time, in milliseconds, that <see cref="Take"/> waits
for a client to be
+ /// returned to the pool before throwing. Defaults to 10000 (10
seconds).
+ /// </summary>
+ public int TimeoutInMs { get; set; } = DefaultTimeoutInMs;
+
+ internal const int DefaultTimeoutInMs = 10_000;
+
+ /// <summary>
+ /// The wait timeout expressed in seconds. Kept for backward
compatibility only; it is a thin
+ /// wrapper over <see cref="TimeoutInMs"/>. Prefer <see
cref="TimeoutInMs"/>, which avoids the
+ /// unit ambiguity that previously caused millisecond values to be
interpreted as seconds.
+ /// </summary>
+ [Obsolete("Use TimeoutInMs instead. This property interprets its value
as seconds.")]
+ public int Timeout
+ {
+ get => TimeoutInMs / 1000;
+ set => TimeoutInMs = value * 1000;
+ }
+
public Client Take()
{
Client client = null;
+ // One overall deadline for the whole call. Return() uses
PulseAll, so every waiter wakes up
+ // while only one of them can dequeue the returned client;
re-arming the full timeout on each
+ // wake-up would let an unlucky waiter exceed the configured bound
indefinitely under churn.
+ var budgetMs = TimeoutInMs;
+ var elapsed = Stopwatch.StartNew();
Monitor.Enter(ClientQueue);
try
{
while (true)
{
- bool timeout = false;
- if (ClientQueue.IsEmpty)
+ if (ClientQueue.TryDequeue(out client))
{
- timeout = !Monitor.Wait(ClientQueue,
TimeSpan.FromSeconds(Timeout));
+ break;
}
- ClientQueue.TryDequeue(out client);
- if (client != null || timeout)
+ var remainingMs = budgetMs -
(int)elapsed.ElapsedMilliseconds;
+ if (remainingMs <= 0)
{
break;
}
+
+ Monitor.Wait(ClientQueue,
TimeSpan.FromMilliseconds(remainingMs));
}
}
finally
@@ -86,7 +112,7 @@ namespace Apache.IoTDB
}
if (client == null)
{
- var reasonPhrase = $"Connection pool is empty and wait time
out({Timeout}s)";
+ var reasonPhrase = $"Connection pool is empty and wait time
out({budgetMs}ms)";
if (DiagnosticReporter != null)
{
throw
DiagnosticReporter.BuildDepletionException(reasonPhrase);
diff --git a/src/Apache.IoTDB/SessionPool.Builder.cs
b/src/Apache.IoTDB/SessionPool.Builder.cs
index 9de2874..7f5dcb4 100644
--- a/src/Apache.IoTDB/SessionPool.Builder.cs
+++ b/src/Apache.IoTDB/SessionPool.Builder.cs
@@ -34,6 +34,7 @@ public partial class SessionPool
private int _poolSize = 8;
private bool _enableRpcCompression = false;
private int _connectionTimeoutInMs = 500;
+ private int _poolWaitTimeoutInMs = DefaultPoolWaitTimeoutInMs;
private bool _useSsl = false;
private string _certificatePath = null;
private string _sqlDialect = IoTDBConstant.TREE_SQL_DIALECT;
@@ -94,6 +95,18 @@ public partial class SessionPool
return this;
}
+ /// <summary>
+ /// Sets how long, in milliseconds, an operation waits for a client to
become available in the pool
+ /// before a <see cref="SessionPoolDepletedException"/> is thrown.
Defaults to
+ /// <see cref="DefaultPoolWaitTimeoutInMs"/> (10 seconds). This is
independent of
+ /// <see cref="SetConnectionTimeoutInMs"/>, which controls the
socket-level timeout.
+ /// </summary>
+ public Builder SetPoolWaitTimeoutInMs(int poolWaitTimeoutInMs)
+ {
+ _poolWaitTimeoutInMs = poolWaitTimeoutInMs;
+ return this;
+ }
+
public Builder SetUseSsl(bool useSsl)
{
_useSsl = useSsl;
@@ -135,6 +148,7 @@ public partial class SessionPool
_poolSize = 8;
_enableRpcCompression = false;
_connectionTimeoutInMs = 500;
+ _poolWaitTimeoutInMs = DefaultPoolWaitTimeoutInMs;
_useSsl = false;
_certificatePath = null;
_sqlDialect = IoTDBConstant.TREE_SQL_DIALECT;
@@ -146,9 +160,9 @@ public partial class SessionPool
// if nodeUrls is not empty, use nodeUrls to create session pool
if (_nodeUrls.Count > 0)
{
- return new SessionPool(_nodeUrls, _username, _password,
_fetchSize, _zoneId, _poolSize, _enableRpcCompression, _connectionTimeoutInMs,
_useSsl, _certificatePath, _sqlDialect, _database);
+ return new SessionPool(_nodeUrls, _username, _password,
_fetchSize, _zoneId, _poolSize, _enableRpcCompression, _connectionTimeoutInMs,
_useSsl, _certificatePath, _sqlDialect, _database, _poolWaitTimeoutInMs);
}
- return new SessionPool(_host, _port, _username, _password,
_fetchSize, _zoneId, _poolSize, _enableRpcCompression, _connectionTimeoutInMs,
_useSsl, _certificatePath, _sqlDialect, _database);
+ return new SessionPool(_host, _port, _username, _password,
_fetchSize, _zoneId, _poolSize, _enableRpcCompression, _connectionTimeoutInMs,
_useSsl, _certificatePath, _sqlDialect, _database, _poolWaitTimeoutInMs);
}
}
}
diff --git a/src/Apache.IoTDB/SessionPool.cs b/src/Apache.IoTDB/SessionPool.cs
index 1c12d85..fff1ea3 100644
--- a/src/Apache.IoTDB/SessionPool.cs
+++ b/src/Apache.IoTDB/SessionPool.cs
@@ -40,6 +40,12 @@ namespace Apache.IoTDB
private static readonly TSProtocolVersion ProtocolVersion =
TSProtocolVersion.IOTDB_SERVICE_PROTOCOL_V3;
private const string DepletionReasonReconnectFailed = "Reconnection
failed";
+ /// <summary>
+ /// Default time, in milliseconds, that an operation waits for a
client to become available
+ /// in the pool before a <see cref="SessionPoolDepletedException"/> is
raised.
+ /// </summary>
+ public const int DefaultPoolWaitTimeoutInMs = 10_000;
+
private readonly string _username;
private readonly string _password;
private bool _enableRpcCompression;
@@ -55,6 +61,11 @@ namespace Apache.IoTDB
/// _timeout is the amount of time a Session will wait for a send
operation to complete successfully.
/// </summary>
private readonly int _timeout;
+ /// <summary>
+ /// _poolWaitTimeoutInMs is the amount of time, in milliseconds, an
operation waits for a client
+ /// to become available in the pool. It is unrelated to the
socket-level _timeout.
+ /// </summary>
+ private readonly int _poolWaitTimeoutInMs = DefaultPoolWaitTimeoutInMs;
private readonly int _poolSize = 4;
private readonly string _sqlDialect = IoTDBConstant.TREE_SQL_DIALECT;
private string _database;
@@ -65,6 +76,14 @@ namespace Apache.IoTDB
private ConcurrentClientQueue _clients;
private ILogger _logger;
private PoolHealthMetrics _healthMetrics;
+ /// <summary>
+ /// Configured capacity that currently holds no connection, because a
reconnection attempt failed and
+ /// the dead connection was discarded. The pool refills this capacity
on demand: a slot is only
+ /// materialized when an acquisition finds no idle client, so under
light or sequential load this
+ /// stays above zero even after the server has fully recovered. It
measures unrealized capacity,
+ /// not server availability.
+ /// </summary>
+ private int _unrealizedCapacity;
public delegate Task<TResult> AsyncOperation<TResult>(Client client);
@@ -83,6 +102,14 @@ namespace Apache.IoTDB
/// </summary>
public int FailedReconnections =>
_healthMetrics?.GetReconnectionFailureTally() ?? 0;
+ /// <summary>
+ /// Configured capacity that currently holds no connection. Capacity
is refilled on demand - a slot is
+ /// materialized only when an acquisition finds no idle client - so a
steady non-zero value under light
+ /// load is normal and does NOT mean the server is unreachable. Use
<see cref="FailedReconnections"/>
+ /// to reason about server availability.
+ /// </summary>
+ public int UnrealizedCapacity => Volatile.Read(ref
_unrealizedCapacity);
+
[Obsolete("This method is deprecated, please use new
SessionPool.Builder().")]
public SessionPool(string host, int port, int poolSize)
@@ -110,6 +137,10 @@ namespace Apache.IoTDB
{
}
protected internal SessionPool(string host, int port, string username,
string password, int fetchSize, string zoneId, int poolSize, bool
enableRpcCompression, int timeout, bool useSsl, string certificatePath, string
sqlDialect, string database)
+ : this(host, port, username, password, fetchSize,
zoneId, poolSize, enableRpcCompression, timeout, useSsl, certificatePath,
sqlDialect, database, DefaultPoolWaitTimeoutInMs)
+ {
+ }
+ protected internal SessionPool(string host, int port, string username,
string password, int fetchSize, string zoneId, int poolSize, bool
enableRpcCompression, int timeout, bool useSsl, string certificatePath, string
sqlDialect, string database, int poolWaitTimeoutInMs)
{
_host = host;
_port = port;
@@ -125,6 +156,7 @@ namespace Apache.IoTDB
_certificatePath = certificatePath;
_sqlDialect = sqlDialect;
_database = database;
+ _poolWaitTimeoutInMs = poolWaitTimeoutInMs;
}
/// <summary>
/// Initializes a new instance of the <see cref="SessionPool"/> class.
@@ -153,6 +185,10 @@ namespace Apache.IoTDB
}
protected internal SessionPool(List<string> nodeUrls, string username,
string password, int fetchSize, string zoneId, int poolSize, bool
enableRpcCompression, int timeout, bool useSsl, string certificatePath, string
sqlDialect, string database)
+ : this(nodeUrls, username, password, fetchSize,
zoneId, poolSize, enableRpcCompression, timeout, useSsl, certificatePath,
sqlDialect, database, DefaultPoolWaitTimeoutInMs)
+ {
+ }
+ protected internal SessionPool(List<string> nodeUrls, string username,
string password, int fetchSize, string zoneId, int poolSize, bool
enableRpcCompression, int timeout, bool useSsl, string certificatePath, string
sqlDialect, string database, int poolWaitTimeoutInMs)
{
if (nodeUrls.Count == 0)
{
@@ -172,10 +208,61 @@ namespace Apache.IoTDB
_certificatePath = certificatePath;
_sqlDialect = sqlDialect;
_database = database;
+ _poolWaitTimeoutInMs = poolWaitTimeoutInMs;
+ }
+ /// <summary>
+ /// Acquires a client from the pool. If the pool has no idle client
but still owns unrealized
+ /// capacity left behind by earlier failed reconnections, one of those
slots is materialized on the
+ /// spot instead of blocking on a queue that nobody will ever feed.
This is what lets the pool
+ /// recover on its own after the server has been unreachable for a
while.
+ /// </summary>
+ private async Task<Client> AcquireClientAsync(CancellationToken
cancellationToken = default)
+ {
+ if (!_isClose && _clients.ClientQueue.IsEmpty &&
TryReserveUnrealizedCapacity())
+ {
+ try
+ {
+ return await Reconnect(cancellationToken:
cancellationToken);
+ }
+ catch (ReconnectionFailedException reconnectEx)
+ {
+ // Still unreachable - hand the slot back so a later call
can retry.
+ Interlocked.Increment(ref _unrealizedCapacity);
+ throw new
SessionPoolDepletedException(DepletionReasonReconnectFailed, AvailableClients,
TotalPoolSize, FailedReconnections, reconnectEx);
+ }
+ catch
+ {
+ // Any other failure must not swallow the slot either.
+ Interlocked.Increment(ref _unrealizedCapacity);
+ throw;
+ }
+ }
+
+ return _clients.Take();
+ }
+
+ /// <summary>
+ /// Atomically claims one unit of unrealized capacity, returning false
when none is left.
+ /// </summary>
+ private bool TryReserveUnrealizedCapacity()
+ {
+ while (true)
+ {
+ int current = Volatile.Read(ref _unrealizedCapacity);
+ if (current <= 0)
+ {
+ return false;
+ }
+ if (Interlocked.CompareExchange(ref _unrealizedCapacity,
current - 1, current) == current)
+ {
+ return true;
+ }
+ }
}
+
public async Task<TResult>
ExecuteClientOperationAsync<TResult>(AsyncOperation<TResult> operation, string
errMsg, bool retryOnFailure = true, bool putClientBack = true)
{
- Client client = _clients.Take();
+ Client client = await AcquireClientAsync();
bool shouldReturnClient = true;
bool operationSucceeded = false;
try
@@ -196,8 +283,11 @@ namespace Apache.IoTDB
}
catch (ReconnectionFailedException reconnectEx)
{
- // Reconnection failed - original client was closed by
Reconnect
+ // Reconnection failed - original client was closed by
Reconnect. Record the now-empty
+ // slot so the pool keeps its configured capacity and
can rebuild it later, instead of
+ // shrinking by one on every failure until every
caller blocks forever.
shouldReturnClient = false;
+ Interlocked.Increment(ref _unrealizedCapacity);
throw new
SessionPoolDepletedException(DepletionReasonReconnectFailed, AvailableClients,
TotalPoolSize, FailedReconnections, reconnectEx);
}
@@ -268,8 +358,9 @@ namespace Apache.IoTDB
{
_healthMetrics = new PoolHealthMetrics(_poolSize);
_clients = new ConcurrentClientQueue();
- _clients.Timeout = _timeout * 5;
+ _clients.TimeoutInMs = _poolWaitTimeoutInMs;
_clients.DiagnosticReporter = this;
+ Volatile.Write(ref _unrealizedCapacity, 0);
if (_nodeUrls.Count == 0)
{
@@ -344,7 +435,9 @@ namespace Apache.IoTDB
}
else
{
- int startIndex = _endPoints.FindIndex(x => x.Ip ==
originalClient.EndPoint.Ip && x.Port == originalClient.EndPoint.Port);
+ int startIndex = originalClient == null
+ ? 0
+ : _endPoints.FindIndex(x => x.Ip ==
originalClient.EndPoint.Ip && x.Port == originalClient.EndPoint.Port);
if (startIndex == -1)
{
throw new ArgumentException($"The original client is not
in the list of endpoints. Original client:
{originalClient.EndPoint.Ip}:{originalClient.EndPoint.Port}");
@@ -372,6 +465,16 @@ namespace Apache.IoTDB
throw new ReconnectionFailedException("Error occurs when
reconnecting session pool. Could not connect to any server");
}
+ /// <summary>
+ /// Indicates whether this pool has been opened and not yet closed by
the caller.
+ /// </summary>
+ /// <remarks>
+ /// This reflects the lifecycle of the pool object only - it is NOT a
server-connectivity probe.
+ /// The client performs no heartbeat, so a server going down does not
flip this back to false;
+ /// it stays true until <see cref="Close"/> is called. To reason about
connectivity, use
+ /// <see cref="AvailableClients"/>, <see cref="UnrealizedCapacity"/>
and <see cref="FailedReconnections"/>,
+ /// or simply let an operation throw <see
cref="SessionPoolDepletedException"/>.
+ /// </remarks>
public bool IsOpen() => !_isClose;
public async Task Close()
@@ -381,6 +484,13 @@ namespace Apache.IoTDB
return;
}
+ // Flip the lifecycle state first. It must not depend on how many
clients happen to be queued:
+ // once every connection has been discarded the queue is empty,
and assigning _isClose
+ // inside the loop below would leave IsOpen() true forever.
Zeroing the unrealized capacity also
+ // disables the rebuild path in AcquireClientAsync while we are
tearing down.
+ _isClose = true;
+ Volatile.Write(ref _unrealizedCapacity, 0);
+
foreach (var client in _clients.ClientQueue.AsEnumerable())
{
var closeSessionRequest = new
TSCloseSessionReq(client.SessionId);
@@ -394,8 +504,6 @@ namespace Apache.IoTDB
}
finally
{
- _isClose = true;
-
client.Transport?.Close();
}
}
@@ -430,7 +538,7 @@ namespace Apache.IoTDB
return _zoneId;
}
- var client = _clients.Take();
+ var client = await AcquireClientAsync();
try
{
diff --git a/src/Apache.IoTDB/TableSessionPool.Builder.cs
b/src/Apache.IoTDB/TableSessionPool.Builder.cs
index 10e24c8..192a677 100644
--- a/src/Apache.IoTDB/TableSessionPool.Builder.cs
+++ b/src/Apache.IoTDB/TableSessionPool.Builder.cs
@@ -37,6 +37,7 @@ public partial class TableSessionPool
private int _poolSize = 8;
private bool _enableRpcCompression = false;
private int _connectionTimeoutInMs = 500;
+ private int _poolWaitTimeoutInMs =
SessionPool.DefaultPoolWaitTimeoutInMs;
private bool _useSsl = false;
private string _certificatePath = null;
private string _sqlDialect = IoTDBConstant.TREE_SQL_DIALECT;
@@ -97,6 +98,18 @@ public partial class TableSessionPool
return this;
}
+ /// <summary>
+ /// Sets how long, in milliseconds, an operation waits for a client to
become available in the pool
+ /// before a <see cref="SessionPoolDepletedException"/> is thrown.
Defaults to
+ /// <see cref="SessionPool.DefaultPoolWaitTimeoutInMs"/> (10 seconds).
This is independent of
+ /// <see cref="SetConnectionTimeoutInMs"/>, which controls the
socket-level timeout.
+ /// </summary>
+ public Builder SetPoolWaitTimeoutInMs(int poolWaitTimeoutInMs)
+ {
+ _poolWaitTimeoutInMs = poolWaitTimeoutInMs;
+ return this;
+ }
+
public Builder SetUseSsl(bool useSsl)
{
_useSsl = useSsl;
@@ -138,6 +151,7 @@ public partial class TableSessionPool
_poolSize = 8;
_enableRpcCompression = false;
_connectionTimeoutInMs = 500;
+ _poolWaitTimeoutInMs = SessionPool.DefaultPoolWaitTimeoutInMs;
_sqlDialect = IoTDBConstant.TABLE_SQL_DIALECT;
_database = "";
}
@@ -148,11 +162,11 @@ public partial class TableSessionPool
// if nodeUrls is not empty, use nodeUrls to create session pool
if (_nodeUrls.Count > 0)
{
- sessionPool = new SessionPool(_nodeUrls, _username, _password,
_fetchSize, _zoneId, _poolSize, _enableRpcCompression, _connectionTimeoutInMs,
_useSsl, _certificatePath, _sqlDialect, _database);
+ sessionPool = new SessionPool(_nodeUrls, _username, _password,
_fetchSize, _zoneId, _poolSize, _enableRpcCompression, _connectionTimeoutInMs,
_useSsl, _certificatePath, _sqlDialect, _database, _poolWaitTimeoutInMs);
}
else
{
- sessionPool = new SessionPool(_host, _port, _username,
_password, _fetchSize, _zoneId, _poolSize, _enableRpcCompression,
_connectionTimeoutInMs, _useSsl, _certificatePath, _sqlDialect, _database);
+ sessionPool = new SessionPool(_host, _port, _username,
_password, _fetchSize, _zoneId, _poolSize, _enableRpcCompression,
_connectionTimeoutInMs, _useSsl, _certificatePath, _sqlDialect, _database,
_poolWaitTimeoutInMs);
}
return new TableSessionPool(sessionPool);
}
diff --git a/tests/Apache.IoTDB.Tests/ConcurrentClientQueueTests.cs
b/tests/Apache.IoTDB.Tests/ConcurrentClientQueueTests.cs
new file mode 100644
index 0000000..8511020
--- /dev/null
+++ b/tests/Apache.IoTDB.Tests/ConcurrentClientQueueTests.cs
@@ -0,0 +1,178 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+using System;
+using System.Diagnostics;
+using System.Threading;
+using System.Threading.Tasks;
+using NUnit.Framework;
+
+namespace Apache.IoTDB.Tests
+{
+ [TestFixture]
+ public class ConcurrentClientQueueTests
+ {
+ private static Client NewStubClient() => new Client(null, 1L, 1L,
null, null);
+
+ [Test]
+ public void TimeoutInMs_DefaultsToTenSeconds()
+ {
+ var queue = new ConcurrentClientQueue();
+
+ Assert.That(queue.TimeoutInMs, Is.EqualTo(10_000), "Default pool
wait timeout should be 10 seconds.");
+ }
+
+ [Test]
+ public void Take_EmptyQueue_ThrowsWithinConfiguredMilliseconds()
+ {
+ // Regression guard: the wait timeout used to be interpreted as
seconds while callers assigned
+ // a millisecond value, so a 200ms budget blocked the caller for
200 seconds instead.
+ var queue = new ConcurrentClientQueue { TimeoutInMs = 200 };
+ var stopwatch = Stopwatch.StartNew();
+
+ Assert.Throws<TimeoutException>(() => queue.Take());
+
+ stopwatch.Stop();
+ Assert.That(stopwatch.ElapsedMilliseconds, Is.LessThan(5_000),
+ "Take() must honour TimeoutInMs as milliseconds, not
seconds.");
+ }
+
+ [Test]
+ public void Take_EmptyQueue_ReportsTimeoutUnitInMessage()
+ {
+ var queue = new ConcurrentClientQueue { TimeoutInMs = 150 };
+
+ var ex = Assert.Throws<TimeoutException>(() => queue.Take());
+
+ Assert.That(ex.Message, Does.Contain("150ms"), "The depletion
message should state the timeout in milliseconds.");
+ }
+
+ [Test]
+ public void Take_ReturnsClientOnceOneIsHandedBack()
+ {
+ var queue = new ConcurrentClientQueue { TimeoutInMs = 5_000 };
+ var expected = NewStubClient();
+
+ var taker = Task.Run(() => queue.Take());
+ Thread.Sleep(100); // let the taker block on the empty queue
+ queue.Return(expected);
+
+ Assert.That(taker.Wait(TimeSpan.FromSeconds(5)), Is.True, "Take()
should be woken by Return().");
+ Assert.That(taker.Result, Is.SameAs(expected));
+ }
+
+ [Test]
+ public void Take_DequeuesWithoutWaitingWhenClientAvailable()
+ {
+ var queue = new ConcurrentClientQueue { TimeoutInMs = 60_000 };
+ var expected = NewStubClient();
+ queue.Add(expected);
+
+ var stopwatch = Stopwatch.StartNew();
+ var actual = queue.Take();
+ stopwatch.Stop();
+
+ Assert.That(actual, Is.SameAs(expected));
+ Assert.That(stopwatch.ElapsedMilliseconds, Is.LessThan(1_000), "A
ready client must be returned immediately.");
+ }
+
+#pragma warning disable CS0618 // exercising the obsolete compatibility shim
on purpose
+ [Test]
+ public void
ObsoleteTimeoutProperty_ConvertsBetweenSecondsAndMilliseconds()
+ {
+ var queue = new ConcurrentClientQueue { Timeout = 3 };
+
+ Assert.That(queue.TimeoutInMs, Is.EqualTo(3_000), "Setting Timeout
(seconds) should scale to milliseconds.");
+ Assert.That(queue.Timeout, Is.EqualTo(3), "Reading Timeout should
scale back to seconds.");
+ }
+#pragma warning restore CS0618
+
+ [Test]
+ public void Take_RepeatedWakeUps_StillHonoursTheOverallDeadline()
+ {
+ // Return() pulses ALL waiters while only one can dequeue, so a
losing waiter re-enters the loop
+ // and waits again. If the full timeout were re-armed on each
wake-up, that waiter could exceed
+ // its budget indefinitely under steady churn. Here the queue is
deliberately never fed: the
+ // waiter is only ever pulsed, so it must still give up once the
overall deadline passes.
+ var queue = new ConcurrentClientQueue { TimeoutInMs = 200 };
+ var stop = new ManualResetEventSlim(false);
+
+ var pulser = Task.Run(() =>
+ {
+ while (!stop.IsSet)
+ {
+ Monitor.Enter(queue.ClientQueue);
+ try
+ {
+ Monitor.PulseAll(queue.ClientQueue);
+ }
+ finally
+ {
+ Monitor.Exit(queue.ClientQueue);
+ }
+ Thread.Sleep(20);
+ }
+ });
+
+ var stopwatch = Stopwatch.StartNew();
+ Assert.Throws<TimeoutException>(() => queue.Take());
+ stopwatch.Stop();
+ stop.Set();
+ pulser.Wait(TimeSpan.FromSeconds(5));
+
+ Assert.That(stopwatch.ElapsedMilliseconds, Is.LessThan(1_000),
+ $"Take() must not re-arm its budget on every wake-up (waited
{stopwatch.ElapsedMilliseconds}ms for a 200ms budget).");
+ }
+
+ [Test]
+ public void Take_MultipleWaiters_EachHonoursTheOverallDeadline()
+ {
+ var queue = new ConcurrentClientQueue { TimeoutInMs = 300 };
+ const int waiterCount = 4;
+
+ var stopwatch = Stopwatch.StartNew();
+ var waiters = new Task[waiterCount];
+ for (int i = 0; i < waiterCount; i++)
+ {
+ waiters[i] = Task.Run(() =>
+ {
+ try
+ {
+ queue.Take();
+ return true; // won the single client
+ }
+ catch (TimeoutException)
+ {
+ return false; // timed out within budget
+ }
+ });
+ }
+
+ Thread.Sleep(100);
+ queue.Return(NewStubClient()); // wakes all waiters; only one can
win
+
+ Assert.That(Task.WaitAll(waiters, TimeSpan.FromSeconds(5)),
Is.True,
+ "Every waiter should settle within its own deadline rather
than waiting indefinitely.");
+ stopwatch.Stop();
+
+ Assert.That(stopwatch.ElapsedMilliseconds, Is.LessThan(2_000),
+ $"Losing waiters must not restart their budget (took
{stopwatch.ElapsedMilliseconds}ms for a 300ms budget).");
+ }
+ }
+}
diff --git a/tests/Apache.IoTDB.Tests/SessionPoolConfigurationTests.cs
b/tests/Apache.IoTDB.Tests/SessionPoolConfigurationTests.cs
new file mode 100644
index 0000000..e27cf96
--- /dev/null
+++ b/tests/Apache.IoTDB.Tests/SessionPoolConfigurationTests.cs
@@ -0,0 +1,135 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+using System.Collections.Generic;
+using System.Reflection;
+using System.Threading.Tasks;
+using NUnit.Framework;
+
+namespace Apache.IoTDB.Tests
+{
+ [TestFixture]
+ public class SessionPoolConfigurationTests
+ {
+ private static int ReadPoolWaitTimeout(SessionPool pool)
+ {
+ var field = typeof(SessionPool).GetField("_poolWaitTimeoutInMs",
BindingFlags.NonPublic | BindingFlags.Instance);
+ Assert.That(field, Is.Not.Null, "_poolWaitTimeoutInMs field is
expected to exist on SessionPool.");
+ return (int)field.GetValue(pool);
+ }
+
+ [Test]
+ public void Builder_WithoutExplicitPoolWaitTimeout_UsesDefault()
+ {
+ var pool = new
SessionPool.Builder().SetHost("127.0.0.1").SetPort(6667).Build();
+
+ Assert.That(ReadPoolWaitTimeout(pool),
Is.EqualTo(SessionPool.DefaultPoolWaitTimeoutInMs));
+ }
+
+ [Test]
+ public void Builder_PoolWaitTimeoutIsIndependentOfConnectionTimeout()
+ {
+ // The pool wait timeout used to be derived from the connection
timeout (_timeout * 5) and then
+ // read back as seconds, which turned a 500ms connection timeout
into a ~41 minute pool wait.
+ var pool = new SessionPool.Builder()
+ .SetHost("127.0.0.1")
+ .SetPort(6667)
+ .SetConnectionTimeoutInMs(500)
+ .SetPoolWaitTimeoutInMs(3_000)
+ .Build();
+
+ Assert.That(ReadPoolWaitTimeout(pool), Is.EqualTo(3_000));
+ }
+
+ [Test]
+ public void Builder_WithNodeUrls_PropagatesPoolWaitTimeout()
+ {
+ var pool = new SessionPool.Builder()
+ .SetNodeUrl(new List<string> { "127.0.0.1:6667",
"127.0.0.1:6668" })
+ .SetPoolWaitTimeoutInMs(7_500)
+ .Build();
+
+ Assert.That(ReadPoolWaitTimeout(pool), Is.EqualTo(7_500));
+ }
+
+ [Test]
+ public void TableSessionPoolBuilder_PropagatesPoolWaitTimeout()
+ {
+ var tablePool = new TableSessionPool.Builder()
+ .SetHost("127.0.0.1")
+ .SetPort(6667)
+ .SetPoolWaitTimeoutInMs(4_200)
+ .Build();
+
+ var innerField = typeof(TableSessionPool).GetField("sessionPool",
BindingFlags.NonPublic | BindingFlags.Instance);
+ Assert.That(innerField, Is.Not.Null, "TableSessionPool is expected
to wrap a SessionPool.");
+
+ var inner = (SessionPool)innerField.GetValue(tablePool);
+ Assert.That(ReadPoolWaitTimeout(inner), Is.EqualTo(4_200));
+ }
+
+ [Test]
+ public void NewPool_StartsWithNoUnrealizedCapacityAndIsNotOpen()
+ {
+ var pool = new
SessionPool.Builder().SetHost("127.0.0.1").SetPort(6667).Build();
+
+ Assert.That(pool.UnrealizedCapacity, Is.Zero);
+ Assert.That(pool.IsOpen(), Is.False);
+ }
+
+ private static void SetPrivateField(object target, string name, object
value)
+ {
+ var field = target.GetType().GetField(name, BindingFlags.NonPublic
| BindingFlags.Instance);
+ Assert.That(field, Is.Not.Null, $"{name} field is expected to
exist.");
+ field.SetValue(target, value);
+ }
+
+ [Test]
+ public async Task Close_EmptyClientQueue_StillMarksThePoolClosed()
+ {
+ // Regression guard: _isClose used to be assigned only inside the
foreach over queued clients.
+ // Once every connection had been discarded the queue was empty,
the loop ran zero times,
+ // and Close() returned while IsOpen() stayed true - leaving the
rebuild path armed.
+ var pool = new
SessionPool.Builder().SetHost("127.0.0.1").SetPort(6667).Build();
+ SetPrivateField(pool, "_clients", new ConcurrentClientQueue());
+ SetPrivateField(pool, "_isClose", false);
+ SetPrivateField(pool, "_unrealizedCapacity", 8);
+
+ Assert.That(pool.IsOpen(), Is.True, "Precondition: the pool looks
open with an empty queue.");
+
+ await pool.Close();
+
+ Assert.That(pool.IsOpen(), Is.False, "Close() must flip the
lifecycle flag regardless of queue contents.");
+ Assert.That(pool.UnrealizedCapacity, Is.Zero, "Close() must disarm
capacity refill.");
+ }
+
+ [Test]
+ public async Task Close_IsIdempotentWhenQueueIsEmpty()
+ {
+ var pool = new
SessionPool.Builder().SetHost("127.0.0.1").SetPort(6667).Build();
+ SetPrivateField(pool, "_clients", new ConcurrentClientQueue());
+ SetPrivateField(pool, "_isClose", false);
+
+ await pool.Close();
+ await pool.Close();
+
+ Assert.That(pool.IsOpen(), Is.False);
+ }
+ }
+}