| | | 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 | | } |