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
129 changes: 85 additions & 44 deletions src/bindings/wasix-ts/tools-package/tools/smoke-host.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@ import { createPackedWasixConsumer } from '../../tools/packed-node-fixture.mjs';
import {
connect,
controlPacket,
expectClosedBeforeReady,
onceClosed,
onceConnected,
readExchange,
Expand All @@ -29,6 +28,8 @@ const verify = await readFile(
const fixture = JSON.parse(
await readFile(resolve(root, 'src/shared/fixtures/postgres/logical-tools.json'), 'utf8'),
);
const socketOperationTimeoutMs = 30_000;
const queuedClientObservationMs = 100;
const scratch = await mkdtemp(join(tmpdir(), `oliphaunt-wasix-${runtimeName}-tools-`));

try {
Expand Down Expand Up @@ -176,80 +177,120 @@ async function verifyServer(openServer, listen) {
const server = await openServer({ listen });
try {
const socket = connect(server.connectionString);
let queued;
let queuedStartup;
try {
await onceConnected(socket);
await withSocketDeadline(socket, onceConnected(socket), 'first client connect');
for (const code of [80_877_103, 80_877_104]) {
const negotiation = withSocketDeadline(
socket,
readSingleByte(socket),
`PostgreSQL negotiation ${code}`,
);
socket.write(controlPacket(code));
const response = await readSingleByte(socket);
const response = await negotiation;
if (response !== 'N'.charCodeAt(0)) {
throw new Error(
`local server returned ${response} for PostgreSQL negotiation request ${code}`,
);
}
}
const firstStartup = withSocketDeadline(socket, readExchange(socket), 'first client startup');
socket.write(startupPacket('postgres', 'postgres'));
await readExchange(socket);
await firstStartup;
console.log(`WASIX TypeScript ${runtimeName} tools/server smoke: first startup`);
const rejected = connect(server.connectionString);
try {
await onceConnected(rejected);
const rejection = expectClosedBeforeReady(rejected);
rejected.write(startupPacket('postgres', 'postgres'));
await rejection;
console.log(
`WASIX TypeScript ${runtimeName} tools/server smoke: concurrent client rejected`,
);
} finally {
rejected.destroy();
}

// The Rust listener deliberately owns one complete client at a time. A
// second TCP/Unix connection can finish its host handshake in the OS
// backlog, but its PostgreSQL startup must wait for the active backend.
queued = connect(server.connectionString);
await withSocketDeadline(queued, onceConnected(queued), 'queued client connect');
queuedStartup = readExchange(queued);
queued.write(startupPacket('postgres', 'postgres'));
await expectStillPending(queuedStartup, queuedClientObservationMs, 'queued client startup');
console.log(`WASIX TypeScript ${runtimeName} tools/server smoke: second client queued`);

const copy = withSocketDeadline(socket, readExchange(socket), 'first client COPY');
socket.write(simpleQuery('COPY (SELECT generate_series(1, 100000)) TO STDOUT'));
const copied = await readExchange(socket);
const copied = await copy;
console.log(`WASIX TypeScript ${runtimeName} tools/server smoke: first COPY`);
if (copied.copyBytes < 500_000) {
throw new Error(`local server truncated COPY output at ${copied.copyBytes} bytes`);
}
const begin = withSocketDeadline(socket, readExchange(socket), 'first client BEGIN');
socket.write(simpleQuery('BEGIN'));
await readExchange(socket);
await begin;
const create = withSocketDeadline(
socket,
readExchange(socket),
'first client transaction query',
);
socket.write(simpleQuery('CREATE TABLE disconnect_must_rollback(value integer)'));
await readExchange(socket);
await create;
const firstClosed = withSocketDeadline(socket, onceClosed(socket), 'first client disconnect');
socket.destroy();
await onceClosed(socket);
await firstClosed;
console.log(`WASIX TypeScript ${runtimeName} tools/server smoke: disconnect recovered`);

const next = await connectWhenReady(server.connectionString);
try {
next.write(simpleQuery('CREATE TABLE disconnect_must_rollback(value integer)'));
await readExchange(next);
next.end(Uint8Array.of('X'.charCodeAt(0), 0, 0, 0, 4));
console.log(`WASIX TypeScript ${runtimeName} tools/server smoke: next client accepted`);
} finally {
next.destroy();
}
await withSocketDeadline(queued, queuedStartup, 'queued client startup after handoff');
const queuedQuery = withSocketDeadline(queued, readExchange(queued), 'queued client query');
queued.write(simpleQuery('CREATE TABLE disconnect_must_rollback(value integer)'));
await queuedQuery;
const queuedClosed = withSocketDeadline(queued, onceClosed(queued), 'queued client close');
queued.end(Uint8Array.of('X'.charCodeAt(0), 0, 0, 0, 4));
await queuedClosed;
console.log(`WASIX TypeScript ${runtimeName} tools/server smoke: queued client accepted`);
} finally {
socket.destroy();
queued?.destroy();
await queuedStartup?.catch(() => undefined);
}
} finally {
await server.close();
}
}

async function connectWhenReady(connectionString) {
const deadline = Date.now() + 30_000;
let lastError;
while (Date.now() < deadline) {
const socket = connect(connectionString);
try {
await onceConnected(socket);
socket.write(startupPacket('postgres', 'postgres'));
await readExchange(socket);
return socket;
} catch (error) {
lastError = error;
function withSocketDeadline(socket, operation, label) {
return new Promise((resolveOperation, rejectOperation) => {
let settled = false;
const finish = (settle, value) => {
if (settled) return;
settled = true;
clearTimeout(timeout);
settle(value);
};
const timeout = setTimeout(() => {
if (settled) return;
settled = true;
socket.destroy();
await new Promise((resolveDelay) => setTimeout(resolveDelay, 25));
}
rejectOperation(new Error(`${label} timed out after ${socketOperationTimeoutMs}ms`));
}, socketOperationTimeoutMs);
operation.then(
(value) => finish(resolveOperation, value),
(error) => finish(rejectOperation, error),
);
});
}

async function expectStillPending(operation, observationMs, label) {
let timeout;
const observed = await Promise.race([
operation.then(
() => ({ status: 'resolved' }),
(error) => ({ status: 'rejected', error }),
),
new Promise((resolveObservation) => {
timeout = setTimeout(() => resolveObservation({ status: 'pending' }), observationMs);
}),
]);
clearTimeout(timeout);
if (observed.status === 'pending') return;
if (observed.status === 'rejected') {
throw new Error(`${label} was rejected instead of waiting behind the active client`, {
cause: observed.error,
});
}
throw new Error(`local server did not accept the next client: ${String(lastError)}`);
throw new Error(`${label} completed while the first client still owned the embedded backend`);
}

function readRuntime(args) {
Expand Down
33 changes: 12 additions & 21 deletions src/bindings/wasix-ts/tools/pgwire-client.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,18 @@ export function onceConnected(socket) {

export function onceClosed(socket) {
if (socket.destroyed) return Promise.resolve();
return new Promise((resolveClosed) => socket.once('close', resolveClosed));
return new Promise((resolveClosed) => {
const onClose = () => {
socket.off('error', onError);
resolveClosed();
};
// A close waiter owns the socket after all protocol operations have
// completed. Keep transport errors from escaping as uncaught EventEmitter
// errors while preserving the helper's close-only completion contract.
const onError = () => {};
socket.once('close', onClose);
socket.on('error', onError);
});
}

export function readSingleByte(socket) {
Expand Down Expand Up @@ -63,26 +74,6 @@ export function readSingleByte(socket) {
});
}

export async function expectClosedBeforeReady(socket) {
try {
await readExchange(socket);
} catch (error) {
if (String(error).includes('closed before ReadyForQuery')) return;
if (
typeof error === 'object' &&
error !== null &&
'code' in error &&
['ECONNRESET', 'EPIPE', 'ERR_SOCKET_CLOSED', 'ERR_STREAM_DESTROYED'].includes(
String(error.code),
)
) {
return;
}
throw error;
}
throw new Error('concurrent local-server client unexpectedly reached ReadyForQuery');
}

export function readExchange(socket) {
return new Promise((resolveExchange, rejectExchange) => {
let buffered = new Uint8Array();
Expand Down
79 changes: 63 additions & 16 deletions src/sdks/rust/tests/native_extensions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@ use std::fs;
use std::future::Future;
use std::path::{Path, PathBuf};
use std::process::Command;
use std::sync::Arc;
use std::sync::{Arc, Mutex};
use std::task::{Context, Poll, Wake, Waker};
use std::thread;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
Expand Down Expand Up @@ -448,14 +448,28 @@ async fn open_extension_database(
root: &Path,
) -> Result<TestDatabase> {
if mode == TestMode::Server {
Ok(TestDatabase::Server(
OliphauntServer::builder()
.storage(DatabaseStorage::Directory(root.to_path_buf()))
.listen(ServerListen::tcp())
.extension(extension)
.start()
.await?,
))
let server = OliphauntServer::builder()
.storage(DatabaseStorage::Directory(root.to_path_buf()))
.listen(ServerListen::tcp())
.extension(extension)
.start()
.await?;
let session = match support::ExternalRawSession::connect(server.connection_string()) {
Ok(session) => session,
Err(error) => {
if let Err(close_error) = server.close().await {
return Err(std::io::Error::other(format!(
"connect native extension server test client: {error}; close server after failed client startup: {close_error}"
))
.into());
}
return Err(error.into());
}
};
Ok(TestDatabase::Server {
server,
session: Mutex::new(session),
})
} else {
Ok(TestDatabase::Embedded(
extension_builder(mode, broker, extension, root)
Expand All @@ -467,31 +481,64 @@ async fn open_extension_database(

enum TestDatabase {
Embedded(Oliphaunt),
Server(OliphauntServer),
Server {
server: OliphauntServer,
session: Mutex<support::ExternalRawSession>,
},
}

impl TestDatabase {
async fn exec_protocol_raw(&self, request: impl AsRef<[u8]>) -> Result<Vec<u8>> {
match self {
Self::Embedded(database) => Ok(database.exec_protocol_raw(request).await?),
Self::Server(server) => Ok(support::external_raw_query(
server.connection_string(),
request,
)?),
Self::Server { session, .. } => Ok(session
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.exec_protocol_raw(request)?),
}
}

async fn backup(&self) -> Result<Vec<u8>> {
match self {
Self::Embedded(database) => Ok(database.backup().await?),
Self::Server(_) => panic!("native server backup must use pg_basebackup"),
Self::Server { .. } => panic!("native server backup must use pg_basebackup"),
}
}

async fn close(&self) -> Result<()> {
match self {
Self::Embedded(database) => Ok(database.close().await?),
Self::Server(database) => Ok(database.close().await?),
Self::Server { server, session } => {
let session_close = session
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.close();
let server_close = server.close().await;
session_close?;
server_close?;
Ok(())
}
}
}
}

impl Drop for TestDatabase {
fn drop(&mut self) {
let close = match self {
Self::Embedded(database) => block_on(database.close()),
Self::Server { server, session } => {
if let Err(error) = session
.get_mut()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.close()
{
eprintln!("close native extension server test client during drop: {error}");
}
block_on(server.close())
}
};
if let Err(error) = close {
eprintln!("close native extension test database during drop: {error}");
}
}
}
Expand Down
Loading
Loading