Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@
import com.google.cloud.spanner.InstanceAdminClient;
import com.google.cloud.spanner.InstanceNotFoundException;
import com.google.cloud.spanner.PGAdapterSessionPoolOptionsHelper;
import com.google.cloud.spanner.SessionPoolOptions;

Check warning on line 31 in src/main/java/com/google/cloud/spanner/pgadapter/ConnectionHandler.java

View workflow job for this annotation

GitHub Actions / e2e-psql-v13-v14

com.google.cloud.spanner.SessionPoolOptions in com.google.cloud.spanner has been deprecated

Check warning on line 31 in src/main/java/com/google/cloud/spanner/pgadapter/ConnectionHandler.java

View workflow job for this annotation

GitHub Actions / e2e-psql-v15-v14

com.google.cloud.spanner.SessionPoolOptions in com.google.cloud.spanner has been deprecated

Check warning on line 31 in src/main/java/com/google/cloud/spanner/pgadapter/ConnectionHandler.java

View workflow job for this annotation

GitHub Actions / e2e-psql-v12-v14

com.google.cloud.spanner.SessionPoolOptions in com.google.cloud.spanner has been deprecated

Check warning on line 31 in src/main/java/com/google/cloud/spanner/pgadapter/ConnectionHandler.java

View workflow job for this annotation

GitHub Actions / e2e-psql-v17-v14

com.google.cloud.spanner.SessionPoolOptions in com.google.cloud.spanner has been deprecated

Check warning on line 31 in src/main/java/com/google/cloud/spanner/pgadapter/ConnectionHandler.java

View workflow job for this annotation

GitHub Actions / e2e-psql-v16-v14

com.google.cloud.spanner.SessionPoolOptions in com.google.cloud.spanner has been deprecated

Check warning on line 31 in src/main/java/com/google/cloud/spanner/pgadapter/ConnectionHandler.java

View workflow job for this annotation

GitHub Actions / e2e-psql-v14-v14

com.google.cloud.spanner.SessionPoolOptions in com.google.cloud.spanner has been deprecated

Check warning on line 31 in src/main/java/com/google/cloud/spanner/pgadapter/ConnectionHandler.java

View workflow job for this annotation

GitHub Actions / integration-test

com.google.cloud.spanner.SessionPoolOptions in com.google.cloud.spanner has been deprecated

Check warning on line 31 in src/main/java/com/google/cloud/spanner/pgadapter/ConnectionHandler.java

View workflow job for this annotation

GitHub Actions / integration-test-emulator (11)

com.google.cloud.spanner.SessionPoolOptions in com.google.cloud.spanner has been deprecated

Check warning on line 31 in src/main/java/com/google/cloud/spanner/pgadapter/ConnectionHandler.java

View workflow job for this annotation

GitHub Actions / ubuntu (11)

com.google.cloud.spanner.SessionPoolOptions in com.google.cloud.spanner has been deprecated

Check warning on line 31 in src/main/java/com/google/cloud/spanner/pgadapter/ConnectionHandler.java

View workflow job for this annotation

GitHub Actions / integration-test-emulator (21)

com.google.cloud.spanner.SessionPoolOptions in com.google.cloud.spanner has been deprecated

Check warning on line 31 in src/main/java/com/google/cloud/spanner/pgadapter/ConnectionHandler.java

View workflow job for this annotation

GitHub Actions / ubuntu (8)

com.google.cloud.spanner.SessionPoolOptions in com.google.cloud.spanner has been deprecated

Check warning on line 31 in src/main/java/com/google/cloud/spanner/pgadapter/ConnectionHandler.java

View workflow job for this annotation

GitHub Actions / ubuntu (21)

com.google.cloud.spanner.SessionPoolOptions in com.google.cloud.spanner has been deprecated

Check warning on line 31 in src/main/java/com/google/cloud/spanner/pgadapter/ConnectionHandler.java

View workflow job for this annotation

GitHub Actions / macos (8)

com.google.cloud.spanner.SessionPoolOptions in com.google.cloud.spanner has been deprecated

Check warning on line 31 in src/main/java/com/google/cloud/spanner/pgadapter/ConnectionHandler.java

View workflow job for this annotation

GitHub Actions / macos (17)

com.google.cloud.spanner.SessionPoolOptions in com.google.cloud.spanner has been deprecated
import com.google.cloud.spanner.Spanner;
import com.google.cloud.spanner.SpannerException;
import com.google.cloud.spanner.SpannerException.ResourceNotFoundException;
Expand Down Expand Up @@ -718,6 +718,15 @@
portal.close();
}

/**
* Unregisters a portal from the connection map without closing it. Returns the unregistered
* portal, or null if no portal was registered with this name. Unlike {@link
* #closePortal(String)}, this method does not throw an exception if the portal is not found.
*/
public IntermediatePortalStatement unregisterPortal(String portalName) {
return this.portalsMap.remove(portalName);
}

public boolean hasPortal(String portalName) {
return this.portalsMap.containsKey(portalName);
}
Expand Down Expand Up @@ -873,6 +882,16 @@
this.statementsMap.remove(statementName);
}

/**
* Unregisters a prepared statement from the connection map without closing it. Returns the
* unregistered statement, or null if no statement was registered with this name. Unlike {@link
* #closeStatement(String)}, this method does not throw an exception if the statement is not
* found.
*/
public IntermediatePreparedStatement unregisterStatement(String statementName) {
return this.statementsMap.remove(statementName);
}

public void closeAllStatements() {
Set<String> names = new HashSet<>(this.statementsMap.keySet());
for (String statementName : names) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,10 @@ public class ExtendedQueryProtocolHandler {
private static final Logger logger =
Logger.getLogger(ExtendedQueryProtocolHandler.class.getName());

private final List<AbstractQueryProtocolMessage> messages = new ArrayList<>();
@VisibleForTesting static final int DEFAULT_BUFFER_CAPACITY = 32;
@VisibleForTesting static final int MAX_BUFFER_CAPACITY = 256;

private final ArrayList<AbstractQueryProtocolMessage> messages;
private final ConnectionHandler connectionHandler;
private final BackendConnection backendConnection;

Expand Down Expand Up @@ -79,9 +82,18 @@ public ExtendedQueryProtocolHandler(ConnectionHandler connectionHandler) {
@VisibleForTesting
public ExtendedQueryProtocolHandler(
ConnectionHandler connectionHandler, BackendConnection backendConnection) {
this(connectionHandler, backendConnection, new ArrayList<>(DEFAULT_BUFFER_CAPACITY));
}

@VisibleForTesting
ExtendedQueryProtocolHandler(
ConnectionHandler connectionHandler,
BackendConnection backendConnection,
ArrayList<AbstractQueryProtocolMessage> messages) {
this.connectionHandler = Preconditions.checkNotNull(connectionHandler);
this.connectionId = connectionHandler.getTraceConnectionId().toString();
this.backendConnection = Preconditions.checkNotNull(backendConnection);
this.messages = Preconditions.checkNotNull(messages);
}

/** Returns the backend PG connection for this query handler. */
Expand Down Expand Up @@ -220,8 +232,20 @@ private void flushMessages(boolean includeReadyResponse) throws Exception {
Logging.format(
"Flushing message", Action.Finished, () -> String.format("Message: %s", message)));
if (message.isReturnedErrorResponse()) {
for (int j = i + 1; j < messages.size(); j++) {
messages.get(j).abort();
// Abort remaining messages in reverse (LIFO) order to correctly unwind state mutations in
// the opposite order that they were registered/buffered (e.g. if a pipeline contains a
// Close followed by a Parse/Bind reusing the same statement or portal name, or both
// creation and closure of a statement in the same aborted pipeline).
for (int j = messages.size() - 1; j > i; j--) {
AbstractQueryProtocolMessage messageToAbort = messages.get(j);
try {
messageToAbort.abort();
} catch (Exception exception) {
logger.log(
Level.WARNING,
exception,
() -> String.format("Failed to abort message: %s", messageToAbort));
}
}
Comment thread
olavloite marked this conversation as resolved.
break;
}
Expand All @@ -239,7 +263,12 @@ private void flushMessages(boolean includeReadyResponse) throws Exception {
throw exception;
} finally {
connectionHandler.getConnectionMetadata().getOutputStream().flush();
boolean shouldTrimToSize = messages.size() > MAX_BUFFER_CAPACITY;
messages.clear();
if (shouldTrimToSize) {
messages.trimToSize();
messages.ensureCapacity(DEFAULT_BUFFER_CAPACITY);
}
logger.log(Level.FINER, Logging.format("Flushing messages", Action.Finished));
endSpan();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,16 +17,20 @@
import com.google.api.core.InternalApi;
import com.google.cloud.spanner.pgadapter.ConnectionHandler;
import com.google.cloud.spanner.pgadapter.error.PGException;
import com.google.cloud.spanner.pgadapter.statements.BackendConnection;
import com.google.cloud.spanner.pgadapter.statements.IntermediatePreparedStatement;
import com.google.cloud.spanner.pgadapter.statements.IntermediateStatement;
import com.google.cloud.spanner.pgadapter.statements.InvalidStatement;
import com.google.cloud.spanner.pgadapter.wireoutput.CloseCompleteResponse;
import java.text.MessageFormat;
import javax.annotation.Nullable;

/** Close the designated statement. */
@InternalApi
public class CloseMessage extends ControlMessage {
public class CloseMessage extends AbstractQueryProtocolMessage {

protected static final char IDENTIFIER = 'C';
private static final String RECEIVED_EVENT_DESCRIPTION = "Received message: '" + IDENTIFIER + "'";

private final PreparedType type;
private final String name;
Expand All @@ -51,19 +55,51 @@ public CloseMessage(ConnectionHandler connection) throws Exception {
this.statement = statement;
}

/** Close the statement server-side and clean up by deleting their metadata locally. */
@Override
protected void sendPayload() throws Exception {
void buffer(BackendConnection backendConnection) throws Exception {
// Unregister the statement or portal immediately from the connection map so that subsequent
// messages in the same pipeline (e.g. a Parse or Bind reusing the same name) will see the name
// available, or attempts to use this closed statement will fail.
// Resource cleanup (statement.close()) is deferred until flush() so that:
// 1. If an earlier message in the pipeline fails, abort() can restore the statement.
// 2. Preceding messages in the pipeline (e.g. Execute) can finish streaming rows before
// closure.
if (this.statement != null) {
if (this.type == PreparedType.Portal) {
this.statement.close(); // Only portals need to be closed server side, since PS is not bound
this.connection.closePortal(this.name);
this.connection.unregisterPortal(this.name);
} else {
this.connection.closeStatement(this.name);
this.connection.unregisterStatement(this.name);
}
}
CloseCompleteResponse.send(this.outputStream);
this.outputStream.flush();
}

@Override
public void flush() throws Exception {
try {
if (this.statement != null) {
this.statement.close();
}
CloseCompleteResponse.send(this.outputStream);
} catch (Exception exception) {
handleError(exception);
}
}

@Override
public void abort() {
// Only restore prepared statements that did not fail during creation (an InvalidStatement
// created by a failed Parse message must not be resurrected).
// Portals are never restored: in PostgreSQL, any error in a pipeline drops all portals.
if (this.type == PreparedType.Statement
&& this.statement instanceof IntermediatePreparedStatement
&& !(this.statement instanceof InvalidStatement)) {
this.connection.registerStatement(this.name, (IntermediatePreparedStatement) this.statement);
}
}

@Override
public String getSql() {
return this.statement == null ? "" : this.statement.getSql();
}

@Override
Expand All @@ -82,6 +118,11 @@ public String getIdentifier() {
return String.valueOf(IDENTIFIER);
}

@Override
public String receivedEventDescription() {
return RECEIVED_EVENT_DESCRIPTION;
}

public String getName() {
return this.name;
}
Expand Down
Loading
Loading