| | 1 | | /* |
| | 2 | | * Copyright 2017 Stanislav Muhametsin. All rights Reserved. |
| | 3 | | * |
| | 4 | | * Licensed under the Apache License, Version 2.0 (the "License"); |
| | 5 | | * you may not use this file except in compliance with the License. |
| | 6 | | * You may obtain a copy of the License at |
| | 7 | | * |
| | 8 | | * http://www.apache.org/licenses/LICENSE-2.0 |
| | 9 | | * |
| | 10 | | * Unless required by applicable law or agreed to in writing, software |
| | 11 | | * distributed under the License is distributed on an "AS IS" BASIS, |
| | 12 | | * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or |
| | 13 | | * implied. |
| | 14 | | * |
| | 15 | | * See the License for the specific language governing permissions and |
| | 16 | | * limitations under the License. |
| | 17 | | */ |
| | 18 | | using CBAM.Abstractions; |
| | 19 | | using CBAM.Abstractions.Implementation; |
| | 20 | | using CBAM.SQL.Implementation; |
| | 21 | | using CBAM.SQL.PostgreSQL; |
| | 22 | | using CBAM.SQL.PostgreSQL.Implementation; |
| | 23 | | using System; |
| | 24 | | using System.Collections.Generic; |
| | 25 | | using System.IO; |
| | 26 | | using System.Linq; |
| | 27 | | using System.Net; |
| | 28 | | using System.Text; |
| | 29 | | using System.Threading; |
| | 30 | | using System.Threading.Tasks; |
| | 31 | | using UtilPack; |
| | 32 | | using MessageIOArgs = System.ValueTuple<CBAM.SQL.PostgreSQL.BackendABIHelper, System.IO.Stream, System.Threading.Cancell |
| | 33 | | using TBoundTypeInfo = System.ValueTuple<System.Type, CBAM.SQL.PostgreSQL.PgSQLTypeFunctionality, CBAM.SQL.PostgreSQL.Pg |
| | 34 | | using FluentCryptography.SASL; |
| | 35 | | using FluentCryptography.SASL.SCRAM; |
| | 36 | | using AsyncEnumeration.Implementation.Enumerable; |
| | 37 | | using AsyncEnumeration.Implementation.Provider; |
| | 38 | |
|
| | 39 | | #if !NETSTANDARD1_0 |
| | 40 | | using System.Net.Sockets; |
| | 41 | | #endif |
| | 42 | |
|
| | 43 | | namespace CBAM.SQL.PostgreSQL.Implementation |
| | 44 | | { |
| | 45 | | using TSASLAuthState = System.ValueTuple<SASLMechanism, SASLCredentialsSCRAMForClient, ResizableArray<Byte>, IEncodin |
| | 46 | | using TStatementExecutionSimpleTaskParameter = System.ValueTuple<SQLStatementExecutionResult, Func<ValueTask<(Boolean |
| | 47 | |
|
| | 48 | | internal sealed partial class PostgreSQLProtocol : SQLConnectionFunctionalitySU<PgSQLConnectionVendorFunctionality> |
| | 49 | | { |
| | 50 | |
|
| | 51 | |
|
| | 52 | | private Int32 _lastSeenTransactionStatus; |
| | 53 | | private readonly IDictionary<String, String> _serverParameters; |
| | 54 | | //private Int32 _standardConformingStrings; |
| | 55 | | private readonly Version _serverVersion; |
| | 56 | |
|
| | 57 | | public PostgreSQLProtocol( |
| | 58 | | PgSQLConnectionVendorFunctionality vendorFunctionality, |
| | 59 | | Boolean disableBinaryProtocolSend, |
| | 60 | | Boolean disableBinaryProtocolReceive, |
| | 61 | | BackendABIHelper messageIOArgs, |
| | 62 | | Stream stream, |
| | 63 | | ResizableArray<Byte> buffer, |
| | 64 | | IDictionary<String, String> serverParameters, |
| | 65 | | TransactionStatus status, |
| | 66 | | Int32 backendPID |
| | 67 | | #if !NETSTANDARD1_0 |
| | 68 | | , Socket socket |
| | 69 | | #endif |
| 24 | 70 | | ) : base( vendorFunctionality, DefaultAsyncProvider.Instance ) |
| | 71 | | { |
| 24 | 72 | | this.DisableBinaryProtocolSend = disableBinaryProtocolSend; |
| 24 | 73 | | this.DisableBinaryProtocolReceive = disableBinaryProtocolReceive; |
| 24 | 74 | | this.MessageIOArgs = ArgumentValidator.ValidateNotNull( nameof( messageIOArgs ), messageIOArgs ); |
| 24 | 75 | | this.Stream = ArgumentValidator.ValidateNotNull( nameof( stream ), stream ); |
| | 76 | | #if !NETSTANDARD1_0 |
| 24 | 77 | | this.Socket = socket; |
| | 78 | | #endif |
| 24 | 79 | | this.Buffer = buffer ?? new ResizableArray<Byte>( 8, exponentialResize: true ); |
| 24 | 80 | | this.DataRowColumnSizes = new ResizableArray<ResettableTransformable<Int32?, Int32>>( exponentialResize: false |
| 24 | 81 | | this._serverParameters = ArgumentValidator.ValidateNotNull( nameof( serverParameters ), serverParameters ); |
| 24 | 82 | | this.ServerParameters = new System.Collections.ObjectModel.ReadOnlyDictionary<String, String>( serverParameters |
| 28 | 83 | | this.TypeRegistry = new TypeRegistryImpl( vendorFunctionality, sql => this.PrepareStatementForExecution( vendor |
| | 84 | |
|
| 24 | 85 | | if ( serverParameters.TryGetValue( "server_version", out var serverVersionString ) ) |
| | 86 | | { |
| | 87 | | // Parse server version |
| 24 | 88 | | var i = 0; |
| 24 | 89 | | var version = serverVersionString.Trim(); |
| 120 | 90 | | while ( i < version.Length && ( Char.IsDigit( version[i] ) || version[i] == '.' ) ) |
| | 91 | | { |
| 96 | 92 | | ++i; |
| | 93 | | } |
| 24 | 94 | | this._serverVersion = new Version( version.Substring( 0, i ) ); |
| | 95 | |
|
| | 96 | | } |
| | 97 | |
|
| | 98 | | // Min supported version is 8.4. |
| 24 | 99 | | var serverVersion = this._serverVersion; |
| 24 | 100 | | if ( serverVersion != null && ( serverVersion.Major < 8 || ( serverVersion.Major == 8 && serverVersion.Minor < |
| | 101 | | { |
| 0 | 102 | | throw new PgSQLException( "Unsupported server version: " + serverVersion + "." ); |
| | 103 | | } |
| 24 | 104 | | this.LastSeenTransactionStatus = status; |
| 24 | 105 | | this.BackendProcessID = backendPID; |
| 24 | 106 | | this.EnqueuedNotifications = new Queue<NotificationEventArgs>(); |
| 24 | 107 | | } |
| | 108 | |
|
| 507 | 109 | | public TypeRegistryImpl TypeRegistry { get; } |
| | 110 | |
|
| 1 | 111 | | public Int32 BackendProcessID { get; } |
| | 112 | |
|
| 24 | 113 | | public IReadOnlyDictionary<String, String> ServerParameters { get; } |
| | 114 | |
|
| | 115 | | protected override ReservedForStatement CreateReservationObject( SQLStatementBuilderInformation stmt ) |
| | 116 | | { |
| 92 | 117 | | return new PgReservedForStatement( |
| 92 | 118 | | #if DEBUG |
| 92 | 119 | | stmt, |
| 92 | 120 | | #endif |
| 92 | 121 | | stmt.IsSimple(), |
| 92 | 122 | | stmt.HasBatchParameters() ? "cbam_statement" : null |
| 92 | 123 | | ); |
| | 124 | | } |
| | 125 | |
|
| | 126 | | protected override void ValidateStatementOrThrow( SQLStatementBuilderInformation statement ) |
| | 127 | | { |
| 93 | 128 | | ArgumentValidator.ValidateNotNull( nameof( statement ), statement ); |
| 97 | 129 | | if ( statement.BatchParameterCount > 1 ) |
| | 130 | | { |
| | 131 | | // Verify that all columns have same typeIDs |
| 2 | 132 | | var first = statement |
| 2 | 133 | | .GetParametersEnumerable( 0 ) |
| 3 | 134 | | .Select( param => this.TypeRegistry.TryGetTypeInfo( param.ParameterCILType ).DatabaseData.TypeID ) |
| 2 | 135 | | .ToArray(); |
| 2 | 136 | | var max = statement.BatchParameterCount; |
| 12 | 137 | | for ( var i = 1; i < max; ++i ) |
| | 138 | | { |
| 4 | 139 | | var j = 0; |
| 12 | 140 | | foreach ( var param in statement.GetParametersEnumerable( i ) ) |
| | 141 | | { |
| 2 | 142 | | if ( first[j] != this.TypeRegistry.TryGetTypeInfo( param.ParameterCILType ).DatabaseData.TypeID ) |
| | 143 | | { |
| 0 | 144 | | throw new PgSQLException( "When using batch parameters, columns must have same type IDs for all bat |
| | 145 | | } |
| 2 | 146 | | ++j; |
| | 147 | | } |
| | 148 | | } |
| | 149 | | } |
| 92 | 150 | | } |
| | 151 | |
|
| | 152 | | private static (Int32[] ParameterIndices, TypeFunctionalityInformation[] TypeInfos, Int32[] TypeIDs) GetVariablesF |
| | 153 | | SQLStatementBuilderInformation stmt, |
| | 154 | | TypeRegistry typeRegistry, |
| | 155 | | Func<SQLStatementBuilderInformation, Int32, StatementParameter> paramExtractor |
| | 156 | | ) |
| | 157 | | { |
| 41 | 158 | | var pCount = stmt.SQLParameterCount; |
| | 159 | | TypeFunctionalityInformation[] typeInfos; |
| | 160 | | Int32[] typeIDs; |
| 44 | 161 | | if ( pCount > 0 ) |
| | 162 | | { |
| 43 | 163 | | typeInfos = new TypeFunctionalityInformation[pCount]; |
| 40 | 164 | | typeIDs = new Int32[pCount]; |
| 185 | 165 | | for ( var i = 0; i < pCount; ++i ) |
| | 166 | | { |
| 50 | 167 | | var param = paramExtractor( stmt, i ); |
| 48 | 168 | | var typeInfo = typeRegistry.TryGetTypeInfo( param.ParameterCILType ); |
| 49 | 169 | | typeInfos[i] = typeInfo; |
| 51 | 170 | | typeIDs[i] = typeInfo?.DatabaseData?.TypeID ?? 0; |
| | 171 | | } |
| 47 | 172 | | } |
| | 173 | | else |
| | 174 | | { |
| 1 | 175 | | typeInfos = Empty<TypeFunctionalityInformation>.Array; |
| 1 | 176 | | typeIDs = Empty<Int32>.Array; |
| | 177 | | } |
| | 178 | |
|
| 48 | 179 | | return (( (PgSQLStatementBuilderInformation) stmt ).ParameterIndices, typeInfos, typeIDs); |
| | 180 | | } |
| | 181 | |
|
| | 182 | | private MessageIOArgs GetIOArgs( ResizableArray<Byte> bufferToUse = null, CancellationToken? tokenToUse = null ) |
| | 183 | | { |
| 774 | 184 | | return (this.MessageIOArgs, this.Stream, tokenToUse ?? this.CurrentCancellationToken, bufferToUse ?? this.Buffe |
| | 185 | | } |
| | 186 | |
|
| | 187 | | protected override async ValueTask<TStatementExecutionSimpleTaskParameter> ExecuteStatementAsBatch( |
| | 188 | | SQLStatementBuilderInformation statement, |
| | 189 | | ReservedForStatement reservedState |
| | 190 | | ) |
| | 191 | | { |
| | 192 | | // TODO somehow make statement name and chunk size parametrizable |
| 3 | 193 | | (var parameterIndices, var typeInfos, var typeIDs) = GetVariablesForExtendedQuerySequence( statement, this.Type |
| 2 | 194 | | var ioArgs = this.GetIOArgs(); |
| 2 | 195 | | var stmtName = ( (PgReservedForStatement) reservedState ).StatementName; |
| 2 | 196 | | var chunkSize = 1000; |
| | 197 | |
|
| | 198 | | // Send a parse message with statement name |
| 2 | 199 | | await new ParseMessage( statement.SQL, parameterIndices, typeIDs, stmtName ).SendMessageAsync( ioArgs, true ); |
| | 200 | |
|
| | 201 | | // Now send describe message |
| 2 | 202 | | await new DescribeMessage( true, stmtName ).SendMessageAsync( ioArgs, true ); |
| | 203 | |
|
| | 204 | | // And then Flush message for backend to send responses |
| 2 | 205 | | await FrontEndMessageWithNoContent.FLUSH.SendMessageAsync( ioArgs, false ); |
| | 206 | |
|
| | 207 | | // Receive first batch of messages |
| 2 | 208 | | BackendMessageObject msg = null; |
| 2 | 209 | | SQLStatementExecutionResult current = null; |
| 2 | 210 | | List<PgSQLError> notices = new List<PgSQLError>(); |
| 2 | 211 | | var sendBatch = true; |
| 8 | 212 | | while ( msg == null ) |
| | 213 | | { |
| 6 | 214 | | msg = ( await this.ReadMessagesUntilMeaningful( notices ) ).Item1; |
| 6 | 215 | | switch ( msg ) |
| | 216 | | { |
| | 217 | | case MessageWithNoContents nc: |
| 4 | 218 | | switch ( nc.Code ) |
| | 219 | | { |
| | 220 | | case BackendMessageCode.ParseComplete: |
| | 221 | | // Continue reading messages |
| 2 | 222 | | msg = null; |
| 2 | 223 | | break; |
| | 224 | | case BackendMessageCode.EmptyQueryResponse: |
| | 225 | | // The statement does not produce any data, we are done |
| 0 | 226 | | sendBatch = false; |
| 0 | 227 | | break; |
| | 228 | | case BackendMessageCode.NoData: |
| | 229 | | // Do nothing, thus causing batch messages to be sent |
| | 230 | | break; |
| | 231 | | default: |
| 0 | 232 | | throw new PgSQLException( "Unrecognized response at this point: " + msg.Code ); |
| | 233 | | } |
| | 234 | | break; |
| | 235 | | case RowDescription rd: |
| | 236 | | // This happens when e.g. doing SELECT schema.function(x, y, z) -> can return NULLs or rows, we don't |
| | 237 | | break; // throw new PgSQLException( "Batch statements may only be used for non-query statements." ); |
| | 238 | | case ParameterDescription pd: |
| 2 | 239 | | if ( !ArrayEqualityComparer<Int32>.ArrayEquality( pd.ObjectIDs, typeIDs ) ) |
| | 240 | | { |
| 0 | 241 | | throw new PgSQLException( "Backend required certain amount of parameters, but either they were not |
| | 242 | | } |
| | 243 | | // Continue to RowDescription/NoData message |
| 2 | 244 | | msg = null; |
| 2 | 245 | | break; |
| | 246 | | default: |
| 0 | 247 | | throw new PgSQLException( "Unrecognized response at this point: " + msg.Code ); |
| | 248 | | } |
| | 249 | | } |
| | 250 | |
|
| 2 | 251 | | if ( sendBatch ) |
| | 252 | | { |
| 2 | 253 | | var batchCount = statement.BatchParameterCount; |
| 2 | 254 | | var affectedRowsArray = new Int32[batchCount]; |
| | 255 | | // Send and receive messages asynchronously |
| 2 | 256 | | var commandTag = new String[1]; |
| 2 | 257 | | await |
| 2 | 258 | | #if NET40 |
| 2 | 259 | | TaskEx |
| 2 | 260 | | #else |
| 2 | 261 | | Task |
| 2 | 262 | | #endif |
| 2 | 263 | | .WhenAll( |
| 2 | 264 | | this.SendMessagesForBatch( statement, typeInfos, stmtName, ioArgs, chunkSize, batchCount ), |
| 2 | 265 | | this.ReceiveMessagesForBatch( notices, affectedRowsArray, commandTag ) |
| 2 | 266 | | ); |
| 2 | 267 | | current = new BatchCommandExecutionResultImpl( |
| 2 | 268 | | commandTag[0], |
| 2 | 269 | | new Lazy<SQLException[]>( () => notices?.Select( n => new PgSQLException( n ) )?.ToArray() ), |
| 2 | 270 | | affectedRowsArray |
| 2 | 271 | | ); |
| 2 | 272 | | } |
| | 273 | |
|
| 2 | 274 | | return (current, null); |
| 2 | 275 | | } |
| | 276 | |
|
| | 277 | | private async Task SendMessagesForBatch( |
| | 278 | | SQLStatementBuilderInformation statement, |
| | 279 | | TypeFunctionalityInformation[] typeInfos, |
| | 280 | | String statementName, |
| | 281 | | MessageIOArgs ioArgs, |
| | 282 | | Int32 chunkSize, |
| | 283 | | Int32 batchCount |
| | 284 | | ) |
| | 285 | | { |
| 2 | 286 | | var singleRowParamCount = statement.SQLParameterCount; |
| | 287 | | Int32 max; |
| 2 | 288 | | var execMessage = new ExecuteMessage(); |
| 8 | 289 | | for ( var i = 0; i < batchCount; i = max ) |
| | 290 | | { |
| 2 | 291 | | max = Math.Min( batchCount, i + chunkSize ); |
| 16 | 292 | | for ( var j = i; j < max; ++j ) |
| | 293 | | { |
| | 294 | | // Send Bind and Execute messages |
| | 295 | | // TODO reuse BindMessage -> add Reset method. |
| 6 | 296 | | await new BindMessage( |
| 6 | 297 | | statement.GetParametersEnumerable( j ), |
| 6 | 298 | | singleRowParamCount, |
| 6 | 299 | | typeInfos, |
| 6 | 300 | | this.DisableBinaryProtocolSend, |
| 6 | 301 | | this.DisableBinaryProtocolReceive, |
| 6 | 302 | | statementName: statementName |
| 6 | 303 | | ).SendMessageAsync( ioArgs, true ); |
| 6 | 304 | | await execMessage.SendMessageAsync( ioArgs, true ); |
| | 305 | | } |
| | 306 | |
|
| | 307 | | // Now send flush message for backend to start sending results back |
| 2 | 308 | | await FrontEndMessageWithNoContent.FLUSH.SendMessageAsync( ioArgs, false ); |
| | 309 | | } |
| 2 | 310 | | } |
| | 311 | |
|
| | 312 | | private async Task ReceiveMessagesForBatch( |
| | 313 | | List<PgSQLError> notices, |
| | 314 | | Int32[] affectedRows, |
| | 315 | | String[] commandTag // This is fugly, but other option is to make both ReceiveMessagesForBatch and SendMessages |
| | 316 | | ) |
| | 317 | | { |
| | 318 | | // We must allocate new buffer, since the reading will be done concurrently while the writing still performs |
| | 319 | | // Furthermore, if some error is occurred during sending task, the backend will send error response right away. |
| 2 | 320 | | var buffer = new ResizableArray<Byte>( initialSize: 8, exponentialResize: true ); |
| | 321 | |
|
| 16 | 322 | | for ( var i = 0; i < affectedRows.Length; ++i ) |
| | 323 | | { |
| 6 | 324 | | var msg = ( await this.ReadMessagesUntilMeaningful( notices, bufferToUse: buffer ) ).Item1; |
| 6 | 325 | | if ( msg is MessageWithNoContents nc && msg.Code == BackendMessageCode.BindComplete ) |
| | 326 | | { |
| | 327 | | // Bind was successul - now read result of execute message |
| 6 | 328 | | msg = null; |
| 12 | 329 | | while ( msg == null ) |
| | 330 | | { |
| | 331 | | Int32 remaining; |
| 6 | 332 | | (msg, remaining) = await this.ReadMessagesUntilMeaningful( notices, bufferToUse: buffer ); |
| 6 | 333 | | switch ( msg ) |
| | 334 | | { |
| | 335 | | case CommandComplete cc: |
| 6 | 336 | | Interlocked.Exchange( ref affectedRows[i], cc.AffectedRows ?? 0 ); |
| 6 | 337 | | if ( commandTag[0] == null ) |
| | 338 | | { |
| 2 | 339 | | Interlocked.Exchange( ref commandTag[0], cc.CommandTag ); |
| | 340 | | } |
| 2 | 341 | | break; |
| | 342 | | case DataRowObject dr: |
| | 343 | | // Skip thru data |
| 0 | 344 | | await this.Stream.ReadSpecificAmountAsync( buffer.SetCapacityAndReturnArray( remaining ), 0, rem |
| | 345 | | // And read more |
| 0 | 346 | | msg = null; |
| 0 | 347 | | break; |
| | 348 | | default: |
| 0 | 349 | | throw new PgSQLException( "Unrecognized response at this point: " + msg.Code ); |
| | 350 | | } |
| | 351 | | } |
| 6 | 352 | | } |
| | 353 | | else |
| | 354 | | { |
| 0 | 355 | | throw new PgSQLException( "Unrecognized response at this point: " + msg.Code ); |
| | 356 | | } |
| | 357 | | } |
| 2 | 358 | | } |
| | 359 | |
|
| | 360 | | protected override async ValueTask<TStatementExecutionSimpleTaskParameter> ExecuteStatementAsPrepared( |
| | 361 | | SQLStatementBuilderInformation statement, |
| | 362 | | ReservedForStatement reservedState |
| | 363 | | ) |
| | 364 | | { |
| 86 | 365 | | (var parameterIndices, var typeInfos, var typeIDs) = GetVariablesForExtendedQuerySequence( statement, this.Type |
| 45 | 366 | | var ioArgs = this.GetIOArgs(); |
| | 367 | |
|
| | 368 | | // First, send the parse message |
| 47 | 369 | | await new ParseMessage( statement.SQL, parameterIndices, typeIDs ).SendMessageAsync( ioArgs, true ); |
| | 370 | |
|
| | 371 | | // Then send bind message |
| 46 | 372 | | var bindMsg = new BindMessage( statement.GetParametersEnumerable(), parameterIndices.Length, typeInfos, this.Di |
| 45 | 373 | | await bindMsg.SendMessageAsync( ioArgs, true ); |
| | 374 | |
|
| | 375 | | // Then send describe message |
| 45 | 376 | | await new DescribeMessage( false ).SendMessageAsync( ioArgs, true ); |
| | 377 | |
|
| | 378 | | // Then execute message |
| 46 | 379 | | await new ExecuteMessage().SendMessageAsync( ioArgs, true ); |
| | 380 | |
|
| | 381 | | // Then flush in order to receive response |
| 47 | 382 | | await FrontEndMessageWithNoContent.FLUSH.SendMessageAsync( ioArgs, false ); |
| | 383 | |
|
| | 384 | | // Start receiving messages |
| 47 | 385 | | BackendMessageObject msg = null; |
| 47 | 386 | | SQLStatementExecutionResult current = null; |
| 47 | 387 | | Func<ValueTask<(Boolean, SQLStatementExecutionResult)>> moveNext = null; |
| 47 | 388 | | RowDescription seenRD = null; |
| 47 | 389 | | List<PgSQLError> notices = new List<PgSQLError>(); |
| 232 | 390 | | while ( msg == null ) |
| | 391 | | { |
| 188 | 392 | | msg = ( await this.ReadMessagesUntilMeaningful( notices ) ).Item1; |
| | 393 | | switch ( msg ) |
| | 394 | | { |
| | 395 | | case MessageWithNoContents nc: |
| 93 | 396 | | switch ( nc.Code ) |
| | 397 | | { |
| | 398 | | case BackendMessageCode.ParseComplete: |
| | 399 | | case BackendMessageCode.BindComplete: |
| | 400 | | case BackendMessageCode.NoData: |
| | 401 | | // Continue reading messages |
| 93 | 402 | | msg = null; |
| 93 | 403 | | break; |
| | 404 | | case BackendMessageCode.EmptyQueryResponse: |
| | 405 | | // The statement does not produce any data, we are done |
| | 406 | | break; |
| | 407 | | default: |
| 0 | 408 | | throw new PgSQLException( "Unrecognized response at this point: " + msg.Code ); |
| | 409 | | } |
| | 410 | | break; |
| | 411 | | case RowDescription rd: |
| | 412 | | // 0..* DataRowObjects incoming... |
| 46 | 413 | | seenRD = rd; |
| 47 | 414 | | msg = null; |
| 47 | 415 | | break; |
| | 416 | | case DataRowObject dr: |
| 47 | 417 | | var streamArray = new PgSQLDataRowColumn[seenRD.Fields.Length]; |
| 47 | 418 | | var mdArray = new PgSQLDataColumnMetaDataImpl[streamArray.Length]; |
| 47 | 419 | | PgSQLDataRowColumn prevCol = null; |
| 186 | 420 | | for ( var i = 0; i < streamArray.Length; ++i ) |
| | 421 | | { |
| 47 | 422 | | var curField = seenRD.Fields[i]; |
| 47 | 423 | | var curMD = new PgSQLDataColumnMetaDataImpl( this, curField.DataFormat, curField.dataTypeID, this.T |
| 47 | 424 | | var curStream = new PgSQLDataRowColumn( curMD, i, prevCol, this, reservedState, curField ); |
| 47 | 425 | | prevCol = curStream; |
| 47 | 426 | | streamArray[i] = curStream; |
| 46 | 427 | | curStream.Reset( dr ); |
| 47 | 428 | | mdArray[i] = curMD; |
| | 429 | | } |
| 46 | 430 | | var warningsLazy = LazyFactory.NewReadOnlyResettableLazy<SQLException[]>( () => notices?.Select( n => |
| 47 | 431 | | var dataRowCurrent = new SQLDataRowImpl( |
| 47 | 432 | | new PgSQLDataRowMetaDataImpl( mdArray ), |
| 47 | 433 | | streamArray, |
| 47 | 434 | | warningsLazy |
| 47 | 435 | | ); |
| 46 | 436 | | current = dataRowCurrent; |
| 54 | 437 | | moveNext = async () => await this.MoveNextAsync( reservedState, streamArray, notices, dataRowCurrent, |
| 46 | 438 | | break; |
| | 439 | | case CommandComplete cc: |
| 0 | 440 | | if ( seenRD == null ) |
| | 441 | | { |
| 0 | 442 | | current = new SingleCommandExecutionResultImpl( |
| 0 | 443 | | cc.CommandTag, |
| 0 | 444 | | new Lazy<SQLException[]>( () => notices?.Select( n => new PgSQLException( n ) )?.ToArray() ), |
| 0 | 445 | | cc.AffectedRows ?? 0 |
| 0 | 446 | | ); |
| | 447 | | } |
| 0 | 448 | | break; |
| | 449 | | default: |
| 0 | 450 | | throw new PgSQLException( "Unrecognized response at this point: " + msg.Code ); |
| | 451 | | } |
| | 452 | | } |
| | 453 | |
|
| 46 | 454 | | return (current, moveNext); |
| 47 | 455 | | } |
| | 456 | |
|
| | 457 | | protected override async ValueTask<TStatementExecutionSimpleTaskParameter> ExecuteStatementAsSimple( |
| | 458 | | SQLStatementBuilderInformation stmt, |
| | 459 | | ReservedForStatement reservedState |
| | 460 | | ) |
| | 461 | | { |
| | 462 | | // Send Query message |
| 49 | 463 | | await new QueryMessage( stmt.SQL ).SendMessageAsync( this.GetIOArgs() ); |
| | 464 | |
|
| | 465 | | // Then wait for appropriate response |
| 49 | 466 | | List<PgSQLError> notices = new List<PgSQLError>(); |
| 49 | 467 | | Func<ValueTask<(Boolean, SQLStatementExecutionResult)>> drMoveNext = null; |
| | 468 | |
|
| | 469 | | // We have to always set moveNext, since we might be executing arbitrary amount of SQL statements in simple Sta |
| 49 | 470 | | Func<ValueTask<(Boolean, SQLStatementExecutionResult)>> moveNext = async () => |
| 49 | 471 | | { |
| 175 | 472 | | SQLStatementExecutionResult current = null; |
| 175 | 473 | | if ( drMoveNext != null ) |
| 49 | 474 | | { |
| 49 | 475 | | // We are iterating over some query result, check that first. |
| 122 | 476 | | var drNext = await drMoveNext(); |
| 122 | 477 | | if ( drNext.Item1 ) |
| 49 | 478 | | { |
| 110 | 479 | | current = drNext.Item2; |
| 110 | 480 | | } |
| 49 | 481 | | else |
| 49 | 482 | | { |
| 61 | 483 | | drMoveNext = null; |
| 49 | 484 | | } |
| 49 | 485 | | } |
| 49 | 486 | |
|
| 175 | 487 | | if ( current == null ) |
| 49 | 488 | | { |
| 114 | 489 | | BackendMessageObject msg = null; |
| 114 | 490 | | RowDescription seenRD = null; |
| 226 | 491 | | while ( msg == null ) |
| 49 | 492 | | { |
| 161 | 493 | | msg = ( await this.ReadMessagesUntilMeaningful( notices ) ).Item1; |
| 49 | 494 | |
|
| 49 | 495 | | switch ( msg ) |
| 49 | 496 | | { |
| 49 | 497 | | case CommandComplete cc: |
| 53 | 498 | | if ( seenRD == null ) |
| 49 | 499 | | { |
| 53 | 500 | | current = new SingleCommandExecutionResultImpl( |
| 53 | 501 | | cc.CommandTag, |
| 53 | 502 | | new Lazy<SQLException[]>( () => notices?.Select( n => new PgSQLException( n ) )?.ToArray() |
| 53 | 503 | | cc.AffectedRows ?? 0 |
| 53 | 504 | | ); |
| 53 | 505 | | } |
| 49 | 506 | | else |
| 49 | 507 | | { |
| 49 | 508 | | // RowDescription followed immediately by CommandComplete -> treat as empty query |
| 49 | 509 | | // Read more |
| 49 | 510 | | msg = null; |
| 49 | 511 | | } |
| 53 | 512 | | seenRD = null; |
| 53 | 513 | | break; |
| 49 | 514 | | case RowDescription rd: |
| 96 | 515 | | seenRD = rd; |
| 49 | 516 | | // Read more (DataRow or CommandComplete) |
| 96 | 517 | | msg = null; |
| 96 | 518 | | break; |
| 49 | 519 | | case DataRowObject dr: |
| 49 | 520 | | // First DataRowObject |
| 96 | 521 | | var streamArray = new PgSQLDataRowColumn[seenRD.Fields.Length]; |
| 96 | 522 | | var mdArray = new PgSQLDataColumnMetaDataImpl[streamArray.Length]; |
| 96 | 523 | | PgSQLDataRowColumn prevCol = null; |
| 263 | 524 | | for ( var i = 0; i < streamArray.Length; ++i ) |
| 49 | 525 | | { |
| 109 | 526 | | var curField = seenRD.Fields[i]; |
| 109 | 527 | | var curMD = new PgSQLDataColumnMetaDataImpl( this, curField.DataFormat, curField.dataTypeID, |
| 109 | 528 | | var curStream = new PgSQLDataRowColumn( curMD, i, prevCol, this, reservedState, curField ); |
| 109 | 529 | | prevCol = curStream; |
| 109 | 530 | | streamArray[i] = curStream; |
| 109 | 531 | | curStream.Reset( dr ); |
| 109 | 532 | | mdArray[i] = curMD; |
| 49 | 533 | | } |
| 96 | 534 | | var warningsLazy = LazyFactory.NewReadOnlyResettableLazy<SQLException[]>( () => notices?.Select( |
| 96 | 535 | | var dataRowCurrent = new SQLDataRowImpl( |
| 96 | 536 | | new PgSQLDataRowMetaDataImpl( mdArray ), |
| 96 | 537 | | streamArray, |
| 96 | 538 | | warningsLazy |
| 96 | 539 | | ); |
| 96 | 540 | | current = dataRowCurrent; |
| 169 | 541 | | drMoveNext = async () => await this.MoveNextAsync( reservedState, streamArray, notices, dataRowC |
| 96 | 542 | | break; |
| 49 | 543 | | case ReadyForQuery rfq: |
| 63 | 544 | | ( (PgReservedForStatement) reservedState ).RFQSeen(); |
| 63 | 545 | | break; |
| 49 | 546 | | default: |
| 49 | 547 | | if ( !ReferenceEquals( MessageWithNoContents.EMPTY_QUERY, msg ) ) |
| 49 | 548 | | { |
| 49 | 549 | | throw new PgSQLException( "Unrecognized response at this point: " + msg.Code ); |
| 49 | 550 | | } |
| 49 | 551 | | // Read more |
| 49 | 552 | | msg = null; |
| 49 | 553 | | break; |
| 49 | 554 | | } |
| 49 | 555 | | } |
| 114 | 556 | | } |
| 49 | 557 | |
|
| 175 | 558 | | return (current != null, current); |
| 175 | 559 | | }; |
| | 560 | |
|
| 49 | 561 | | var firstResult = await moveNext(); |
| | 562 | |
|
| 49 | 563 | | return (firstResult.Item1 ? firstResult.Item2 : null, moveNext); |
| | 564 | |
|
| 49 | 565 | | } |
| | 566 | |
|
| | 567 | | private async Task<(Boolean, SQLStatementExecutionResult)> MoveNextAsync( |
| | 568 | | ReservedForStatement reservationObject, |
| | 569 | | PgSQLDataRowColumn[] streams, |
| | 570 | | List<PgSQLError> notices, |
| | 571 | | SQLDataRowImpl dataRow, |
| | 572 | | ReadOnlyResettableLazy<SQLException[]> warningsLazy |
| | 573 | | ) |
| | 574 | | { |
| 80 | 575 | | return await this.UseStreamWithinStatementAsync( reservationObject, async () => |
| 80 | 576 | | { |
| 80 | 577 | | // Force read of all columns |
| 748 | 578 | | foreach ( var colStream in streams ) |
| 80 | 579 | | { |
| 334 | 580 | | await colStream.SkipBytesAsync( this.Buffer.Array ); |
| 80 | 581 | | } |
| 80 | 582 | |
|
| 160 | 583 | | notices.Clear(); |
| 160 | 584 | | var msg = ( await this.ReadMessagesUntilMeaningful( notices ) ).Item1; |
| 160 | 585 | | var dr = msg as DataRowObject; |
| 748 | 586 | | foreach ( var stream in streams ) |
| 80 | 587 | | { |
| 334 | 588 | | stream.Reset( dr ); |
| 80 | 589 | | } |
| 80 | 590 | |
|
| 160 | 591 | | var retVal = dr != null; |
| 160 | 592 | | warningsLazy.Reset(); |
| 160 | 593 | | return (Success: retVal, Item: dataRow); |
| 160 | 594 | | } ); |
| 80 | 595 | | } |
| | 596 | |
|
| | 597 | |
|
| | 598 | | public TransactionStatus LastSeenTransactionStatus |
| | 599 | | { |
| | 600 | | get |
| | 601 | | { |
| 0 | 602 | | return (TransactionStatus) this._lastSeenTransactionStatus; |
| | 603 | | } |
| | 604 | | private set |
| | 605 | | { |
| 122 | 606 | | Interlocked.Exchange( ref this._lastSeenTransactionStatus, (Int32) value ); |
| 122 | 607 | | } |
| | 608 | | } |
| | 609 | |
|
| | 610 | | //public Boolean StandardConformingStrings |
| | 611 | | //{ |
| | 612 | | // get |
| | 613 | | // { |
| | 614 | | // return Convert.ToBoolean( this._standardConformingStrings ); |
| | 615 | | // } |
| | 616 | | // set |
| | 617 | | // { |
| | 618 | | // Interlocked.Exchange( ref this._standardConformingStrings, Convert.ToInt32( value ) ); |
| | 619 | | // } |
| | 620 | | //} |
| | 621 | |
|
| | 622 | | protected override async Task PerformDisposeStatementAsync( |
| | 623 | | ReservedForStatement reservationObject |
| | 624 | | ) |
| | 625 | | { |
| 93 | 626 | | var ioArgs = this.GetIOArgs(); |
| 99 | 627 | | var pgReserved = (PgReservedForStatement) reservationObject; |
| 99 | 628 | | if ( !String.IsNullOrEmpty( pgReserved.StatementName ) ) |
| | 629 | | { |
| | 630 | | // Need to close our named statement |
| 2 | 631 | | await new CloseMessage( true, pgReserved.StatementName ).SendMessageAsync( ioArgs, true ); |
| | 632 | | } |
| | 633 | |
|
| | 634 | | // Simple statement already received RFQ in its MoveNext method |
| 95 | 635 | | if ( !pgReserved.IsSimple ) |
| | 636 | | { |
| | 637 | | // Need to send SYNC |
| 48 | 638 | | await FrontEndMessageWithNoContent.SYNC.SendMessageAsync( ioArgs ); |
| | 639 | |
|
| | 640 | | } |
| | 641 | |
|
| | 642 | | // TODO The new moveNextEnded parameter could tell that instead of RFQEncountered property, investigate that |
| 99 | 643 | | if ( !pgReserved.RFQEncountered ) |
| | 644 | | { |
| | 645 | | // Then wait for RFQ |
| | 646 | | // This happens for non-simple statements, or simple statements which cause exception when iterated over. |
| | 647 | | BackendMessageObject msg; |
| | 648 | | Int32 remaining; |
| 165 | 649 | | while ( ( (msg, remaining) = ( await this.ReadMessagesUntilMeaningful( null, dontThrowExceptions: true ) ) ) |
| | 650 | | { |
| 80 | 651 | | if ( remaining > 0 ) |
| | 652 | | { |
| 0 | 653 | | ioArgs.Item4.CurrentMaxCapacity = remaining; |
| 0 | 654 | | await ioArgs.Item2.ReadSpecificAmountAsync( ioArgs.Item4.Array, 0, remaining, ioArgs.Item3 ); |
| | 655 | | } |
| | 656 | | } |
| | 657 | | } |
| 99 | 658 | | } |
| | 659 | |
|
| 1436 | 660 | | public BackendABIHelper MessageIOArgs { get; } |
| | 661 | |
|
| 1881 | 662 | | public ResizableArray<Byte> Buffer { get; } |
| | 663 | |
|
| 1471 | 664 | | public Stream Stream { get; } |
| | 665 | |
|
| | 666 | | #if !NETSTANDARD1_0 |
| 6 | 667 | | public Socket Socket { get; } |
| | 668 | | #endif |
| | 669 | |
|
| 564 | 670 | | public ResizableArray<ResettableTransformable<Int32?, Int32>> DataRowColumnSizes { get; } |
| | 671 | |
|
| 46 | 672 | | public Boolean DisableBinaryProtocolSend { get; } |
| 49 | 673 | | public Boolean DisableBinaryProtocolReceive { get; } |
| | 674 | |
|
| 4 | 675 | | public Queue<NotificationEventArgs> EnqueuedNotifications { get; } |
| | 676 | |
|
| | 677 | | internal async ValueTask<Object> ConvertFromBytes( |
| | 678 | | Int32 typeID, |
| | 679 | | DataFormat dataFormat, |
| | 680 | | EitherOr<ReservedForStatement, Stream> stream, |
| | 681 | | Int32 byteCount |
| | 682 | | ) |
| | 683 | | { |
| 328 | 684 | | var actualStream = stream.IsFirst ? this.Stream : stream.Second; |
| 328 | 685 | | var typeInfo = this.TypeRegistry.TryGetTypeInfo( typeID ); |
| 329 | 686 | | if ( typeInfo != null ) |
| | 687 | | { |
| 137 | 688 | | var limitedStream = StreamFactory.CreateLimitedReader( |
| 137 | 689 | | actualStream, |
| 137 | 690 | | byteCount, |
| 137 | 691 | | this.CurrentCancellationToken, |
| 137 | 692 | | this.Buffer |
| 137 | 693 | | ); |
| | 694 | |
|
| | 695 | | try |
| | 696 | | { |
| 137 | 697 | | return await typeInfo.Functionality.ReadBackendValueAsync( |
| 137 | 698 | | dataFormat, |
| 137 | 699 | | typeInfo.DatabaseData, |
| 137 | 700 | | this.MessageIOArgs, |
| 137 | 701 | | limitedStream |
| 137 | 702 | | ); |
| | 703 | | } |
| | 704 | | finally |
| | 705 | | { |
| | 706 | | try |
| | 707 | | { |
| 136 | 708 | | await limitedStream.SkipThroughRemainingBytes(); |
| 137 | 709 | | } |
| 0 | 710 | | catch |
| | 711 | | { |
| | 712 | | // Ignore this one. |
| 0 | 713 | | } |
| | 714 | |
|
| | 715 | | } |
| | 716 | |
|
| 0 | 717 | | } |
| 192 | 718 | | else if ( dataFormat == DataFormat.Text ) |
| | 719 | | { |
| | 720 | | // Initial type load, or unknown type and format is textual |
| 192 | 721 | | await actualStream.ReadSpecificAmountAsync( this.Buffer, 0, byteCount, this.CurrentCancellationToken ); |
| 192 | 722 | | return this.MessageIOArgs.GetStringWithPool( this.Buffer.Array, 0, byteCount ); |
| | 723 | | } |
| | 724 | | else |
| | 725 | | { |
| | 726 | | // Unknown type, and data format is binary. |
| 0 | 727 | | throw new PgSQLException( $"The type ID {typeID} is not known." ); |
| | 728 | | } |
| 327 | 729 | | } |
| | 730 | |
|
| | 731 | | internal async ValueTask<(BackendMessageObject, Int32)> ReadMessagesUntilMeaningful( |
| | 732 | | List<PgSQLError> notices, |
| | 733 | | Func<Boolean> checkReadForNextMessage = null, |
| | 734 | | ResizableArray<Byte> bufferToUse = null, |
| | 735 | | Boolean dontThrowExceptions = false |
| | 736 | | ) |
| | 737 | | { |
| | 738 | | Boolean encounteredMeaningful; |
| 561 | 739 | | var ioArgs = this.GetIOArgs( bufferToUse ); |
| | 740 | | BackendMessageObject msg; |
| | 741 | | Int32 remaining; |
| | 742 | | do |
| | 743 | | { |
| 564 | 744 | | (msg, remaining) = await BackendMessageObject.ReadBackendMessageAsync( ioArgs, this.DataRowColumnSizes ); |
| 563 | 745 | | switch ( msg ) |
| | 746 | | { |
| | 747 | | case PgSQLErrorObject errorObject: |
| 0 | 748 | | encounteredMeaningful = false; |
| 0 | 749 | | if ( errorObject.Code == BackendMessageCode.NoticeResponse ) |
| | 750 | | { |
| 0 | 751 | | if ( notices != null ) |
| | 752 | | { |
| 0 | 753 | | notices.Add( ( (PgSQLErrorObject) msg ).Error ); |
| | 754 | | } |
| 0 | 755 | | } |
| 0 | 756 | | else if ( !dontThrowExceptions ) |
| | 757 | | { |
| 0 | 758 | | throw new PgSQLException( ( (PgSQLErrorObject) msg ).Error ); |
| | 759 | | } |
| | 760 | | break; |
| | 761 | | case NotificationMessage notification: |
| 1 | 762 | | this.EnqueuedNotifications.Enqueue( notification.Args ); |
| 1 | 763 | | encounteredMeaningful = false; |
| 1 | 764 | | break; |
| | 765 | | case ParameterStatus ps: |
| 0 | 766 | | this._serverParameters[ps.Name] = ps.Value; |
| 0 | 767 | | encounteredMeaningful = false; |
| 0 | 768 | | break; |
| | 769 | | default: |
| | 770 | | { |
| 561 | 771 | | if ( msg is ReadyForQuery rfq ) |
| | 772 | | { |
| 98 | 773 | | this.LastSeenTransactionStatus = rfq.Status; |
| | 774 | | } |
| 562 | 775 | | encounteredMeaningful = true; |
| | 776 | | break; |
| | 777 | | } |
| | 778 | |
|
| | 779 | | } |
| 563 | 780 | | } while ( !encounteredMeaningful && ( checkReadForNextMessage?.Invoke() ?? true ) ); |
| | 781 | |
|
| 563 | 782 | | return (msg, remaining); |
| 564 | 783 | | } |
| | 784 | |
|
| | 785 | | public async Task PerformClose( CancellationToken token ) |
| | 786 | | { |
| | 787 | | // Send termination message |
| | 788 | | // Don't use this.CurrentCancellationToken, since one-time pool has already reset the token. |
| | 789 | | // Furthermore, we might come here from other entrypoints than connection pool's UseConnection (e.g. when dispo |
| 24 | 790 | | await FrontEndMessageWithNoContent.TERMINATION.SendMessageAsync( this.GetIOArgs( tokenToUse: token ) ); |
| 24 | 791 | | } |
| | 792 | |
|
| | 793 | | #if !NETSTANDARD1_0 |
| | 794 | | private Boolean SocketHasDataPending() |
| | 795 | | { |
| 3 | 796 | | var socket = this.Socket; |
| 3 | 797 | | return socket.Available > 0 || socket.Poll( 1, SelectMode.SelectRead ) || socket.Available > 0; |
| | 798 | | } |
| | 799 | | #endif |
| | 800 | |
|
| | 801 | | public async ValueTask<NotificationEventArgs[]> CheckNotificationsAsync() |
| | 802 | | { |
| | 803 | | // TODO this could be optimized a little, if we notice EnqueuedNotifications.Count > 0, then just don't read fr |
| 2 | 804 | | NotificationEventArgs[] args = null; |
| | 805 | |
|
| | 806 | | NotificationEventArgs[] GetEnqueuedNotifications() |
| | 807 | | { |
| 0 | 808 | | var enqueued = this.EnqueuedNotifications.ToArray(); |
| 0 | 809 | | this.EnqueuedNotifications.Clear(); |
| 0 | 810 | | return enqueued; |
| | 811 | | } |
| | 812 | |
|
| | 813 | | #if !NETSTANDARD1_0 |
| 2 | 814 | | var socket = this.Socket; |
| 2 | 815 | | if ( socket == null ) |
| | 816 | | { |
| | 817 | | #endif |
| | 818 | | // Just do "SELECT 1"; to get any notifications |
| 0 | 819 | | var enumerable = this.PrepareStatementForExecution( this.VendorFunctionality.CreateStatementBuilder( "SELECT |
| 0 | 820 | | .AsObservable(); |
| | 821 | | // Use GetEnqueuedNotifications while we are still inside statement reservation region, by registering to Be |
| 0 | 822 | | enumerable.BeforeEnumerationEnd += ( eArgs ) => args = GetEnqueuedNotifications(); |
| 0 | 823 | | await enumerable.EnumerateAsync(); |
| | 824 | | #if !NETSTANDARD1_0 |
| 0 | 825 | | } |
| | 826 | | else |
| | 827 | | { |
| | 828 | | // First, check from the socket that we have any data pending |
| | 829 | |
|
| 2 | 830 | | var hasDataPending = this.SocketHasDataPending(); |
| 2 | 831 | | if ( hasDataPending || this.EnqueuedNotifications.Count > 0 ) |
| | 832 | | { |
| | 833 | | // There is pending data |
| | 834 | | // We always must use UseStreamOutsideStatementAsync method, since modifying this.EnqueuedNotifications o |
| 0 | 835 | | await this.UseStreamOutsideStatementAsync( async () => |
| 0 | 836 | | { |
| 0 | 837 | | // If we call "ReadMessagesUntilMeaningful" with no socket data pending, we will never break free of l |
| 0 | 838 | | if ( hasDataPending ) |
| 0 | 839 | | { |
| 0 | 840 | | await this.ReadMessagesUntilMeaningful( |
| 0 | 841 | | null, |
| 0 | 842 | | this.SocketHasDataPending |
| 0 | 843 | | ); |
| 0 | 844 | | } |
| 0 | 845 | | args = GetEnqueuedNotifications(); |
| 0 | 846 | | return false; |
| 0 | 847 | | } ); |
| | 848 | | } |
| | 849 | | } |
| | 850 | | #endif |
| | 851 | |
|
| 2 | 852 | | return args ?? Empty<NotificationEventArgs>.Array; |
| | 853 | |
|
| 2 | 854 | | } |
| | 855 | |
|
| | 856 | |
|
| | 857 | | public IAsyncEnumerable<NotificationEventArgs> ListenToNotificationsAsync() |
| | 858 | | { |
| | 859 | | #if !NETSTANDARD1_0 |
| 1 | 860 | | if ( this.Socket == null ) |
| | 861 | | { |
| | 862 | | #else |
| | 863 | | throw new NotSupportedException( "No socket available for this method." ); |
| | 864 | | #endif |
| | 865 | | #if !NETSTANDARD1_0 |
| | 866 | | } |
| | 867 | |
|
| 1 | 868 | | var enqueued = this.EnqueuedNotifications; |
| | 869 | | Boolean KeepReadingMore() |
| | 870 | | { |
| 1 | 871 | | return enqueued.Count <= 0 || ( enqueued.Count <= 1000 && this.SocketHasDataPending() ); |
| | 872 | | } |
| | 873 | |
|
| | 874 | | async Task PerformReadForNotifications() |
| | 875 | | { |
| 1 | 876 | | if ( enqueued.Count <= 0 ) |
| | 877 | | { |
| 1 | 878 | | await this.ReadMessagesUntilMeaningful( null, KeepReadingMore ); |
| | 879 | | } |
| 1 | 880 | | } |
| | 881 | |
|
| 1 | 882 | | return AsyncEnumerationFactory.CreateStatefulWrappingEnumerable( () => |
| 1 | 883 | | { |
| 2 | 884 | | PgReservedForStatement reservation = null; |
| 2 | 885 | | return AsyncEnumerationFactory.CreateWrappingStartInfo( |
| 2 | 886 | | async () => |
| 2 | 887 | | { |
| 3 | 888 | | if ( reservation == null ) |
| 2 | 889 | | { |
| 3 | 890 | | reservation = new PgReservedForStatement( |
| 3 | 891 | | #if DEBUG |
| 3 | 892 | | null, |
| 3 | 893 | | #endif |
| 3 | 894 | | true, |
| 3 | 895 | | null |
| 3 | 896 | | ); |
| 3 | 897 | | reservation.RFQSeen(); |
| 3 | 898 | | await this.UseStreamOutsideStatementAsync( reservation, PerformReadForNotifications, false, true ); |
| 3 | 899 | | } |
| 2 | 900 | | else |
| 2 | 901 | | { |
| 2 | 902 | | await this.UseStreamWithinStatementAsync( reservation, PerformReadForNotifications, true ); |
| 2 | 903 | | } |
| 2 | 904 | |
|
| 3 | 905 | | return enqueued.Count > 0; |
| 3 | 906 | | }, |
| 2 | 907 | | ( out Boolean success ) => |
| 2 | 908 | | { |
| 3 | 909 | | success = enqueued.Count > 0; |
| 3 | 910 | | return success ? enqueued.Dequeue() : default; |
| 2 | 911 | | }, |
| 2 | 912 | | () => |
| 2 | 913 | | { |
| 3 | 914 | | return this.DisposeStatementAsync( reservation ); |
| 2 | 915 | | } |
| 2 | 916 | | ); |
| 1 | 917 | | }, this.AsyncProvider ); |
| | 918 | | #endif |
| | 919 | | } |
| | 920 | |
|
| | 921 | | public static async Task<(PostgreSQLProtocol Protocol, List<PgSQLError> notices)> PerformStartup( |
| | 922 | | PgSQLConnectionVendorFunctionality vendorFunctionality, |
| | 923 | | PgSQLConnectionCreationInfo creationInfo, |
| | 924 | | CancellationToken token, |
| | 925 | | Stream stream, |
| | 926 | | BackendABIHelper abiHelper, |
| | 927 | | ResizableArray<Byte> buffer |
| | 928 | | #if !NETSTANDARD1_0 |
| | 929 | | , Socket socket |
| | 930 | | #endif |
| | 931 | | ) |
| | 932 | | { |
| 24 | 933 | | var initData = creationInfo?.CreationData?.Initialization ?? throw new PgSQLException( "Please specify initiali |
| 24 | 934 | | var startupInfo = await DoConnectionInitialization( |
| 24 | 935 | | creationInfo, |
| 24 | 936 | | (abiHelper, stream, token, buffer) |
| 24 | 937 | | ); |
| 24 | 938 | | var protoConfig = initData?.Protocol; |
| 24 | 939 | | var retVal = ( |
| 24 | 940 | | new PostgreSQLProtocol( |
| 24 | 941 | | vendorFunctionality, |
| 24 | 942 | | protoConfig?.DisableBinaryProtocolSend ?? false, |
| 24 | 943 | | protoConfig?.DisableBinaryProtocolReceive ?? false, |
| 24 | 944 | | abiHelper, |
| 24 | 945 | | stream, |
| 24 | 946 | | buffer, |
| 24 | 947 | | startupInfo.ServerParameters, |
| 24 | 948 | | startupInfo.TransactionStatus, |
| 24 | 949 | | startupInfo.backendProcessID ?? 0 |
| 24 | 950 | | #if !NETSTANDARD1_0 |
| 24 | 951 | | , socket |
| 24 | 952 | | #endif |
| 24 | 953 | | ), |
| 24 | 954 | | startupInfo.Notices ?? new List<PgSQLError>() |
| 24 | 955 | | ); |
| | 956 | |
|
| 24 | 957 | | await retVal.Item1.ReadTypesFromServer( protoConfig?.ForceTypeIDLoad ?? false, token ); |
| | 958 | |
|
| 24 | 959 | | return retVal; |
| 24 | 960 | | } |
| | 961 | |
|
| | 962 | | internal const String SERVER_PARAMETER_DATABASE = "database"; |
| | 963 | |
|
| | 964 | | private static async Task<(IDictionary<String, String> ServerParameters, Int32? backendProcessID, Int32? backendKe |
| | 965 | | PgSQLConnectionCreationInfo creationInfo, |
| | 966 | | MessageIOArgs ioArgs |
| | 967 | | ) |
| | 968 | | { |
| 24 | 969 | | var dbConfig = creationInfo?.CreationData?.Initialization?.Database ?? throw new ArgumentException( "Please spe |
| 24 | 970 | | var authConfig = creationInfo?.CreationData?.Initialization?.Authentication ?? throw new ArgumentException( "Pl |
| | 971 | |
|
| 24 | 972 | | var encoding = ioArgs.Item1.Encoding.Encoding; |
| 24 | 973 | | var username = authConfig.Username ?? throw new ArgumentException( "Please specify username in authentication c |
| 24 | 974 | | var parameters = new Dictionary<String, String>() |
| 24 | 975 | | { |
| 24 | 976 | | { SERVER_PARAMETER_DATABASE, dbConfig.Name ?? throw new ArgumentException("Please specify database name in d |
| 24 | 977 | | { "user",username }, |
| 24 | 978 | | { "DateStyle", "ISO" }, |
| 24 | 979 | | { "client_encoding", encoding.WebName }, |
| 24 | 980 | | { "extra_float_digits", "2" }, |
| 24 | 981 | | { "lc_monetary", "C" } |
| 24 | 982 | | }; |
| 24 | 983 | | var sp = dbConfig.SearchPath; |
| 24 | 984 | | if ( !String.IsNullOrEmpty( sp ) ) |
| | 985 | | { |
| 0 | 986 | | parameters.Add( "search_path", sp ); |
| | 987 | | } |
| | 988 | |
|
| 24 | 989 | | await new StartupMessage( 3 << 16, parameters ).SendMessageAsync( ioArgs ); |
| | 990 | |
|
| | 991 | | BackendMessageObject msg; |
| 24 | 992 | | List<PgSQLError> notices = null; |
| 24 | 993 | | Int32? backendProcessID = null; |
| 24 | 994 | | Int32? backendKeyData = null; |
| 24 | 995 | | TransactionStatus tStatus = 0; |
| 24 | 996 | | Object saslState = null; |
| | 997 | | try |
| | 998 | | { |
| | 999 | | do |
| | 1000 | | { |
| | 1001 | | Int32 ignored; |
| 364 | 1002 | | (msg, ignored) = await BackendMessageObject.ReadBackendMessageAsync( ioArgs, null ); |
| 364 | 1003 | | switch ( msg ) |
| | 1004 | | { |
| | 1005 | | case ParameterStatus ps: |
| 264 | 1006 | | parameters[ps.Name] = ps.Value; |
| 264 | 1007 | | break; |
| | 1008 | | case AuthenticationResponse auth: |
| 52 | 1009 | | var newSaslState = await ProcessAuth( |
| 52 | 1010 | | creationInfo, |
| 52 | 1011 | | username, |
| 52 | 1012 | | ioArgs, |
| 52 | 1013 | | auth, |
| 52 | 1014 | | saslState |
| 52 | 1015 | | ); |
| 52 | 1016 | | if ( newSaslState != null ) |
| | 1017 | | { |
| 8 | 1018 | | saslState = newSaslState; |
| | 1019 | | } |
| 8 | 1020 | | break; |
| | 1021 | | case PgSQLErrorObject error: |
| 0 | 1022 | | if ( error.Code == BackendMessageCode.NoticeResponse ) |
| | 1023 | | { |
| 0 | 1024 | | if ( notices == null ) |
| | 1025 | | { |
| 0 | 1026 | | notices = new List<PgSQLError>(); |
| | 1027 | | } |
| 0 | 1028 | | notices.Add( error.Error ); |
| 0 | 1029 | | } |
| | 1030 | | else |
| | 1031 | | { |
| 0 | 1032 | | throw new PgSQLException( error.Error ); |
| | 1033 | | } |
| | 1034 | | break; |
| | 1035 | | case BackendKeyData key: |
| 24 | 1036 | | backendProcessID = key.ProcessID; |
| 24 | 1037 | | backendKeyData = key.Key; |
| 24 | 1038 | | break; |
| | 1039 | | case ReadyForQuery rfq: |
| 24 | 1040 | | tStatus = rfq.Status; |
| | 1041 | | break; |
| | 1042 | | } |
| 364 | 1043 | | } while ( msg.Code != BackendMessageCode.ReadyForQuery ); |
| 24 | 1044 | | } |
| | 1045 | | finally |
| | 1046 | | { |
| 24 | 1047 | | DisposeSASLState( saslState ); |
| | 1048 | | } |
| 24 | 1049 | | return (parameters, backendProcessID, backendKeyData, notices, tStatus); |
| 24 | 1050 | | } |
| | 1051 | |
|
| | 1052 | | private static async Task<Object> ProcessAuth( |
| | 1053 | | PgSQLConnectionCreationInfo creationInfo, |
| | 1054 | | String username, |
| | 1055 | | MessageIOArgs ioArgs, |
| | 1056 | | AuthenticationResponse msg, |
| | 1057 | | Object saslState |
| | 1058 | | ) |
| | 1059 | | { |
| 52 | 1060 | | var authType = msg.RequestType; |
| 52 | 1061 | | var initData = creationInfo.CreationData.Initialization.Database; |
| 52 | 1062 | | switch ( authType ) |
| | 1063 | | { |
| | 1064 | | case AuthenticationResponse.AuthenticationRequestType.AuthenticationClearTextPassword: |
| 0 | 1065 | | await new PasswordMessage( GetPasswordBytes( creationInfo, ioArgs ) ).SendMessageAsync( ioArgs ); |
| 0 | 1066 | | break; |
| | 1067 | | case AuthenticationResponse.AuthenticationRequestType.AuthenticationMD5Password: |
| 22 | 1068 | | await HandleMD5Authentication( ioArgs, msg, username, GetPasswordBytes( creationInfo, ioArgs ) ).SendMess |
| 22 | 1069 | | break; |
| | 1070 | | case AuthenticationResponse.AuthenticationRequestType.AuthenticationOk: |
| | 1071 | | // Nothing to do |
| | 1072 | | break; |
| | 1073 | | case AuthenticationResponse.AuthenticationRequestType.AuthenticationSASL: |
| 2 | 1074 | | var saslResult = ( HandleSASLAuthentication_Start( creationInfo, ioArgs, username, msg ) ); |
| 2 | 1075 | | saslState = saslResult.Item2; |
| 2 | 1076 | | await ( saslResult.Item1 ?? throw new PgSQLException( "Authentication failed." ) ).SendMessageAsync( ioAr |
| 2 | 1077 | | break; |
| | 1078 | | case AuthenticationResponse.AuthenticationRequestType.AuthenticationSASLContinue: |
| 2 | 1079 | | await ( HandleSASLAuthentication_Continue( ioArgs, msg, saslState ) ?? throw new PgSQLException( "Authent |
| 2 | 1080 | | break; |
| | 1081 | | case AuthenticationResponse.AuthenticationRequestType.AuthenticationSASLFinal: |
| 2 | 1082 | | HandleSASLAuthentication_Final( creationInfo, ioArgs, msg, saslState ); |
| 2 | 1083 | | break; |
| | 1084 | | default: |
| 0 | 1085 | | throw new PgSQLException( $"Authentication kind {authType} is not support." ); |
| | 1086 | | } |
| | 1087 | |
|
| 52 | 1088 | | return saslState; |
| 52 | 1089 | | } |
| | 1090 | |
|
| | 1091 | | private static Byte[] GetPasswordBytes( |
| | 1092 | | PgSQLConnectionCreationInfo creationInfo, |
| | 1093 | | MessageIOArgs ioArgs |
| | 1094 | | ) |
| | 1095 | | { |
| 22 | 1096 | | var authConfig = creationInfo.CreationData.Initialization.Authentication; |
| 22 | 1097 | | var encoding = ioArgs.Item1.Encoding.Encoding; |
| 22 | 1098 | | return ( String.Equals( PgSQLAuthenticationConfiguration.PasswordByteEncoding.WebName, encoding.WebName ) ? |
| 22 | 1099 | | authConfig.PasswordBytes : |
| 22 | 1100 | | encoding.GetBytes( authConfig.Password ) ) ?? throw new PgSQLException( "Backend requested password, but it |
| | 1101 | | } |
| | 1102 | |
|
| | 1103 | | // Having this in separate method also won't force load of UtilPack.Cryptography assemblies if other than MD5/SASL |
| | 1104 | | private static PasswordMessage HandleMD5Authentication( |
| | 1105 | | MessageIOArgs ioArgs, |
| | 1106 | | AuthenticationResponse msg, |
| | 1107 | | String username, |
| | 1108 | | Byte[] pw |
| | 1109 | | ) |
| | 1110 | | { |
| 22 | 1111 | | var buffer = ioArgs.Item4; |
| 22 | 1112 | | var helper = ioArgs.Item1; |
| | 1113 | |
|
| 22 | 1114 | | if ( pw == null ) |
| | 1115 | | { |
| 0 | 1116 | | throw new PgSQLException( "Backend requested password, but it was not supplied." ); |
| | 1117 | | } |
| 22 | 1118 | | using ( var md5 = new FluentCryptography.Digest.MD5() ) |
| | 1119 | | { |
| | 1120 | | // Extract server salt before using args.Buffer |
| | 1121 | |
|
| 22 | 1122 | | var serverSalt = buffer.Array.CreateArrayCopy( msg.AdditionalDataInfo.offset, msg.AdditionalDataInfo.count ) |
| | 1123 | |
|
| | 1124 | | // Hash password with username as salt |
| 22 | 1125 | | var prehashLength = helper.Encoding.Encoding.GetByteCount( username ) + pw.Length; |
| 22 | 1126 | | buffer.CurrentMaxCapacity = prehashLength; |
| 22 | 1127 | | var idx = 0; |
| 22 | 1128 | | pw.CopyTo( buffer.Array, ref idx, 0, pw.Length ); |
| 22 | 1129 | | helper.Encoding.Encoding.GetBytes( username, 0, username.Length, buffer.Array, pw.Length ); |
| 22 | 1130 | | var hash = md5.ComputeDigest( buffer.Array, 0, prehashLength ); |
| | 1131 | |
|
| | 1132 | | // Write hash as hexadecimal string |
| 22 | 1133 | | buffer.CurrentMaxCapacity = hash.Length * 2 * helper.Encoding.BytesPerASCIICharacter; |
| 22 | 1134 | | idx = 0; |
| 748 | 1135 | | foreach ( var hashByte in hash ) |
| | 1136 | | { |
| 352 | 1137 | | helper.Encoding.WriteHexDecimal( buffer.Array, ref idx, hashByte ); |
| | 1138 | | } |
| | 1139 | |
|
| | 1140 | | // Hash result again with server-provided salt |
| 22 | 1141 | | buffer.CurrentMaxCapacity += serverSalt.Length; |
| 22 | 1142 | | var dummy = 0; |
| 22 | 1143 | | serverSalt.CopyTo( buffer.Array, ref dummy, idx, serverSalt.Length ); |
| 22 | 1144 | | hash = md5.ComputeDigest( buffer.Array, 0, idx + serverSalt.Length ); |
| | 1145 | |
|
| | 1146 | | // Send back string "md5" followed by hexadecimal hash value |
| 22 | 1147 | | buffer.CurrentMaxCapacity = 3 * helper.Encoding.BytesPerASCIICharacter + hash.Length * 2 * helper.Encoding.B |
| 22 | 1148 | | idx = 0; |
| 22 | 1149 | | var array = buffer.Array; |
| 22 | 1150 | | helper.Encoding |
| 22 | 1151 | | .WriteASCIIByte( array, ref idx, (Byte) 'm' ) |
| 22 | 1152 | | .WriteASCIIByte( array, ref idx, (Byte) 'd' ) |
| 22 | 1153 | | .WriteASCIIByte( array, ref idx, (Byte) '5' ); |
| 748 | 1154 | | foreach ( var hashByte in hash ) |
| | 1155 | | { |
| 352 | 1156 | | helper.Encoding.WriteHexDecimal( array, ref idx, hashByte ); |
| | 1157 | | } |
| | 1158 | |
|
| 22 | 1159 | | var retValArray = new Byte[idx + 1]; // Remember string-terminating zero |
| 22 | 1160 | | dummy = 0; |
| 22 | 1161 | | array.CopyTo( retValArray, ref dummy, 0, idx ); |
| 22 | 1162 | | return new PasswordMessage( retValArray ); |
| | 1163 | | } |
| | 1164 | |
|
| | 1165 | |
|
| 22 | 1166 | | } |
| | 1167 | |
|
| | 1168 | | private static (PasswordMessage, Object) HandleSASLAuthentication_Start( |
| | 1169 | | PgSQLConnectionCreationInfo creationInfo, |
| | 1170 | | MessageIOArgs ioArgs, |
| | 1171 | | String username, |
| | 1172 | | AuthenticationResponse msg |
| | 1173 | | ) |
| | 1174 | | { |
| 2 | 1175 | | var idx = msg.AdditionalDataInfo.offset; |
| 2 | 1176 | | var count = msg.AdditionalDataInfo.count; |
| 2 | 1177 | | var buffer = ioArgs.Item4; |
| 6 | 1178 | | while ( count > 0 && buffer.Array[idx + count - 1] == 0 ) |
| | 1179 | | { |
| 4 | 1180 | | --count; |
| | 1181 | | } |
| 2 | 1182 | | var protocolEncoding = ioArgs.Item1.Encoding; |
| 2 | 1183 | | var authSchemes = protocolEncoding.Encoding.GetString( buffer.Array, idx, count ); |
| | 1184 | |
|
| 2 | 1185 | | var mechanismInfo = creationInfo.CreateSASLMechanism?.Invoke( authSchemes ) ?? throw new PgSQLException( "Faile |
| 2 | 1186 | | var mechanism = mechanismInfo.Item1 ?? throw new PgSQLException( "Failed to provide SASL mechanism." ); |
| 2 | 1187 | | var mechanismName = mechanismInfo.Item2 ?? throw new PgSQLException( "Failed to provide SASL mechanism name." ) |
| 2 | 1188 | | var authConfig = creationInfo.CreationData.Initialization.Authentication; |
| 2 | 1189 | | var pwDigest = authConfig.PasswordDigest; |
| 2 | 1190 | | var credentials = pwDigest.IsNullOrEmpty() ? |
| 2 | 1191 | | new SASLCredentialsSCRAMForClient( username, authConfig.Password ) : |
| 2 | 1192 | | new SASLCredentialsSCRAMForClient( username, pwDigest ); |
| 2 | 1193 | | var writeBuffer = new ResizableArray<Byte>(); |
| 2 | 1194 | | var saslEncoding = new UTF8Encoding( false, true ).CreateDefaultEncodingInfo(); |
| 2 | 1195 | | var challengeResult = mechanism.ChallengeAsync( credentials.CreateChallengeArguments( |
| 2 | 1196 | | Empty<Byte>.Array, |
| 2 | 1197 | | -1, |
| 2 | 1198 | | -1, |
| 2 | 1199 | | writeBuffer, |
| 2 | 1200 | | 0, |
| 2 | 1201 | | saslEncoding |
| 2 | 1202 | | ) ).GetResultForceSynchronous(); |
| | 1203 | |
|
| | 1204 | | (PasswordMessage, Object) retVal; |
| 2 | 1205 | | if ( !challengeResult.IsFirst || challengeResult.First.Item2 != SASLChallengeResult.MoreToCome ) |
| | 1206 | | { |
| 0 | 1207 | | retVal = default; |
| 0 | 1208 | | } |
| | 1209 | | else |
| | 1210 | | { |
| | 1211 | | // SASL initial response is: null-terminated string for mechanism name, length of initial response, and init |
| 2 | 1212 | | var bytesWritten = challengeResult.First.Item1; |
| 2 | 1213 | | var pwArray = new Byte[ |
| 2 | 1214 | | protocolEncoding.Encoding.GetByteCount( mechanismName ) + protocolEncoding.BytesPerASCIICharacter |
| 2 | 1215 | | + sizeof( Int32 ) |
| 2 | 1216 | | + bytesWritten |
| 2 | 1217 | | ]; |
| 2 | 1218 | | idx = protocolEncoding.Encoding.GetBytes( mechanismName, 0, mechanismName.Length, pwArray, 0 ) + 1; |
| 2 | 1219 | | pwArray.WritePgInt32( ref idx, bytesWritten ); |
| 2 | 1220 | | var dummy = 0; |
| 2 | 1221 | | writeBuffer.Array.CopyTo( pwArray, ref dummy, idx, bytesWritten ); |
| | 1222 | |
|
| 2 | 1223 | | retVal = ( |
| 2 | 1224 | | new PasswordMessage( pwArray ), |
| 2 | 1225 | | new TSASLAuthState( mechanism, credentials, writeBuffer, saslEncoding ) |
| 2 | 1226 | | ); |
| | 1227 | | } |
| | 1228 | |
|
| 2 | 1229 | | return retVal; |
| | 1230 | |
|
| | 1231 | | } |
| | 1232 | |
|
| | 1233 | | private static PasswordMessage HandleSASLAuthentication_Continue( |
| | 1234 | | MessageIOArgs ioArgs, |
| | 1235 | | AuthenticationResponse msg, |
| | 1236 | | Object state |
| | 1237 | | ) |
| | 1238 | | { |
| 2 | 1239 | | var challengeResult = HandleSASLAuthentication_ContinueOrFinal( ioArgs, msg, state ); |
| | 1240 | | PasswordMessage retVal; |
| 2 | 1241 | | if ( challengeResult.IsSecond || challengeResult.First.Item2 != SASLChallengeResult.MoreToCome ) |
| | 1242 | | { |
| 0 | 1243 | | retVal = default; |
| 0 | 1244 | | } |
| | 1245 | | else |
| | 1246 | | { |
| | 1247 | | // Responses are password messages with whole SASL message as content |
| 2 | 1248 | | retVal = new PasswordMessage( ( (TSASLAuthState) state ).Item3.Array.CreateArrayCopy( 0, challengeResult.Fir |
| | 1249 | | } |
| | 1250 | |
|
| 2 | 1251 | | return retVal; |
| | 1252 | | } |
| | 1253 | |
|
| | 1254 | | private static void HandleSASLAuthentication_Final( |
| | 1255 | | PgSQLConnectionCreationInfo creationInfo, |
| | 1256 | | MessageIOArgs ioArgs, |
| | 1257 | | AuthenticationResponse msg, |
| | 1258 | | Object state |
| | 1259 | | ) |
| | 1260 | | { |
| 2 | 1261 | | var challengeResult = HandleSASLAuthentication_ContinueOrFinal( ioArgs, msg, state ); |
| 2 | 1262 | | if ( challengeResult.IsSecond || challengeResult.First.Item2 != SASLChallengeResult.Completed ) |
| | 1263 | | { |
| 0 | 1264 | | throw new PgSQLException( "Authentication failed." ); |
| | 1265 | | } |
| | 1266 | | else |
| | 1267 | | { |
| | 1268 | | try |
| | 1269 | | { |
| 2 | 1270 | | creationInfo.OnSASLSCRAMSuccess?.Invoke( ( (TSASLAuthState) state ).Item2.PasswordDigest ); |
| 2 | 1271 | | } |
| 0 | 1272 | | catch |
| | 1273 | | { |
| | 1274 | | // Ignore... |
| 0 | 1275 | | } |
| | 1276 | |
|
| 2 | 1277 | | DisposeSASLState( state ); |
| | 1278 | | } |
| 2 | 1279 | | } |
| | 1280 | |
|
| | 1281 | | private static EitherOr<(Int32, SASLChallengeResult), Int32> HandleSASLAuthentication_ContinueOrFinal( |
| | 1282 | | MessageIOArgs ioArgs, |
| | 1283 | | AuthenticationResponse msg, |
| | 1284 | | Object state |
| | 1285 | | ) |
| | 1286 | | { |
| 4 | 1287 | | var idx = msg.AdditionalDataInfo.offset; |
| 4 | 1288 | | var count = msg.AdditionalDataInfo.count; |
| 4 | 1289 | | var buffer = ioArgs.Item4; |
| | 1290 | |
|
| 4 | 1291 | | var saslState = (TSASLAuthState) state; |
| 4 | 1292 | | return saslState.Item1.ChallengeAsync( saslState.Item2.CreateChallengeArguments( |
| 4 | 1293 | | buffer.Array, |
| 4 | 1294 | | idx, |
| 4 | 1295 | | count, |
| 4 | 1296 | | saslState.Item3, |
| 4 | 1297 | | 0, |
| 4 | 1298 | | saslState.Item4 |
| 4 | 1299 | | ) ).GetResultForceSynchronous(); |
| | 1300 | | } |
| | 1301 | |
|
| | 1302 | | private static void DisposeSASLState( Object state ) |
| | 1303 | | { |
| 26 | 1304 | | if ( state is TSASLAuthState saslState ) |
| | 1305 | | { |
| 4 | 1306 | | saslState.Item1?.DisposeSafely(); |
| 4 | 1307 | | saslState.Item3?.Array?.Clear(); |
| | 1308 | | } |
| 26 | 1309 | | } |
| | 1310 | |
|
| | 1311 | | internal class PgReservedForStatement : ReservedForStatement |
| | 1312 | | { |
| | 1313 | | private Int32 _rfqEncountered; |
| | 1314 | |
|
| 94 | 1315 | | public PgReservedForStatement( |
| 94 | 1316 | | #if DEBUG |
| 94 | 1317 | | Object statement, |
| 94 | 1318 | | #endif |
| 94 | 1319 | | Boolean isSimple, |
| 94 | 1320 | | String statementName |
| 94 | 1321 | | ) |
| | 1322 | | #if DEBUG |
| | 1323 | | : base( statement ) |
| | 1324 | | #endif |
| | 1325 | | { |
| 96 | 1326 | | this.IsSimple = isSimple; |
| 98 | 1327 | | this.StatementName = statementName; |
| 96 | 1328 | | this._rfqEncountered = Convert.ToInt32( false ); |
| 97 | 1329 | | } |
| | 1330 | |
|
| 96 | 1331 | | public Boolean IsSimple { get; } |
| | 1332 | |
|
| 97 | 1333 | | public String StatementName { get; } |
| | 1334 | |
|
| 99 | 1335 | | public Boolean RFQEncountered => Convert.ToBoolean( this._rfqEncountered ); |
| | 1336 | |
|
| | 1337 | | public void RFQSeen() |
| | 1338 | | { |
| 15 | 1339 | | Interlocked.Exchange( ref this._rfqEncountered, Convert.ToInt32( true ) ); |
| 15 | 1340 | | } |
| | 1341 | | } |
| | 1342 | |
|
| | 1343 | | } |
| | 1344 | |
|
| | 1345 | | // TODO move to utilpack |
| | 1346 | | internal static class E_TODO |
| | 1347 | | { |
| | 1348 | | public static T GetResultForceSynchronous<T>( this ValueTask<T> task ) |
| | 1349 | | { |
| | 1350 | | return task.IsCompleted ? task.Result : throw new InvalidOperationException( "ValueTask is not completed when i |
| | 1351 | | } |
| | 1352 | | } |
| | 1353 | | } |