mahbot 0.7.3

An autonomous agentic engineering system that manages software development through role separation, subagents, and deterministic diagnostics.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
//! IPC query endpoint for the `mahbot debug` read-only CLI.
//!
//! After `multiprocess_wal` was removed, the debug CLI can no longer open a
//! second physical instance of a live store's database (single-process mode
//! holds the flock). Instead the running instance exposes a local IPC endpoint
//! (a Unix-domain socket on Unix, a named pipe on Windows) that accepts
//! read-only SQL queries and returns the rows. The debug CLI connects to it
//! instead of opening the database directly.
//!
//! Protocol (length-prefixed JSON over a local socket stream):
//! - request  → `[u32 LE payload_len][payload]`, payload is JSON `QueryRequest`.
//! - response → `[u32 LE payload_len][payload]`, payload is JSON `QueryResponse`.
//!
//! Write enforcement is the `PRAGMA query_only=1` guard in
//! [`crate::db::Connection::query_readonly`], which runs and resets the pragma
//! (back to `0`) under a single hold of the connection mutex.

use std::path::{Path, PathBuf};
use std::time::Duration;

use interprocess::local_socket::tokio::prelude::*;
use interprocess::local_socket::{GenericFilePath, ListenerOptions, Name};
use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt};
use tracing::{debug, warn};

use crate::db::{ReadonlyRows, Value};

/// Maximum number of result rows served per IPC query (LIMIT+1 semantics).
pub(crate) const IPC_ROW_LIMIT: usize = 10_000;

/// Socket file name for the debug IPC endpoint (a filesystem socket on Unix).
#[cfg(not(windows))]
const IPC_SOCKET_FILE_NAME: &str = "mahbot-debug.sock";
/// Prefix of the Windows named-pipe name. Named pipes live in a global
/// `\\.\pipe\` namespace with no per-directory scope, and the `GenericFilePath`
/// name type only accepts `\\.\pipe\`-prefixed paths there — the location
/// digest is appended below.
#[cfg(windows)]
const IPC_PIPE_PREFIX: &str = r"\\.\pipe\mahbot-debug-";

/// Total wall-clock bound for the instance-up-but-channel-not-bound retry
/// (lock held before the listener binds during boot / self-update hand-off).
const IPC_BOUND_TIMEOUT_SECS: u64 = 15;

/// [`IPC_BOUND_TIMEOUT_SECS`], overridable with `MAHBOT_IPC_BOUND_TIMEOUT_SECS`
/// (0 = give up after the first attempt) — the same env-override pattern the
/// other bounded waits use, and what lets a test exercise the refusal without
/// waiting out the default.
fn ipc_bound_timeout() -> std::time::Duration {
    crate::util::env_duration_secs("MAHBOT_IPC_BOUND_TIMEOUT_SECS", IPC_BOUND_TIMEOUT_SECS)
}

/// Retry backoff schedule (ms) for the IPC socket-not-yet-bound window.
const IPC_RETRY_BACKOFF_MS: [u64; 6] = [50, 100, 200, 400, 800, 1600];

/// Length (bytes) of the `u32 LE` length prefix on each IPC frame.
const FRAME_LEN: usize = std::mem::size_of::<u32>();

/// Upper bound (bytes) on a single IPC frame payload, enforced on read so a
/// same-user misbehaving client cannot trigger a multi-GiB allocation and OOM
/// the instance. Requests are a few KB; responses are bounded by
/// [`IPC_ROW_LIMIT`], so this is far beyond legitimate use.
const MAX_FRAME_LEN: usize = 64 * 1024 * 1024;

/// Resolve the IPC socket name.
///
/// On Unix this is a filesystem socket file under the storage root. On Windows
/// it is a named pipe scoped to that storage location: the pipe namespace is
/// global, so a machine-wide name would let one instance reach another
/// instance's data. Both the listener and the clients call this so they agree
/// on the endpoint name.
#[must_use]
fn socket_path(storage_root: &Path) -> PathBuf {
    #[cfg(windows)]
    {
        use sha2::{Digest, Sha256};

        // A digest of the location, not the path itself: a pipe name may not
        // contain a backslash, must stay well under the 256-char limit, and is
        // case-insensitive. `canonicalize` gives the real spelling of an
        // existing directory (both sides derive the root from the same
        // `default_config_dir`), falling back to the given path when that
        // directory does not exist yet — no instance can be running there;
        // lowercasing removes case-only differences, and SHA-256 truncated to
        // 16 hex chars keeps the name short and collision-free.
        let resolved =
            std::fs::canonicalize(storage_root).unwrap_or_else(|_| storage_root.to_path_buf());
        let digest = crate::util::hex_string(&Sha256::digest(
            resolved.to_string_lossy().to_lowercase().as_bytes(),
        ));
        PathBuf::from(format!("{IPC_PIPE_PREFIX}{}", &digest[..16]))
    }
    #[cfg(not(windows))]
    {
        storage_root.join(IPC_SOCKET_FILE_NAME)
    }
}

#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)]
#[serde(tag = "t", content = "v", rename_all = "snake_case")]
pub(crate) enum WireValue {
    Null,
    Integer(i64),
    /// serde_json serializes non-finite floats as `null`; `de_real` maps it
    /// back to `NaN` so the IPC path round-trips them like the direct
    /// (no-instance) read path renders them ("NaN").
    Real(#[serde(deserialize_with = "de_real")] f64),
    Text(String),
    /// Base64-encoded blob.
    Blob(String),
}

/// Deserialize a `WireValue::Real` payload: a finite JSON number, or `null`
/// (which serde_json emits for NaN/Inf on the wire) mapped to `NaN`.
fn de_real<'de, D>(d: D) -> std::result::Result<f64, D::Error>
where
    D: serde::Deserializer<'de>,
{
    let v = <Option<f64> as serde::Deserialize>::deserialize(d)?;
    Ok(v.unwrap_or(f64::NAN))
}

/// Decode a base64-encoded blob, defaulting to empty on malformed input.
fn decode_blob(b: &str) -> Vec<u8> {
    crate::util::base64_decode(b).unwrap_or_default()
}

impl WireValue {
    pub(crate) fn from_turso(v: &Value) -> Self {
        match v {
            Value::Null => WireValue::Null,
            Value::Integer(i) => WireValue::Integer(*i),
            Value::Real(f) => WireValue::Real(*f),
            Value::Text(s) => WireValue::Text(s.clone()),
            Value::Blob(b) => WireValue::Blob(crate::util::base64_encode(b)),
        }
    }

    pub(crate) fn to_turso(&self) -> Value {
        match self {
            WireValue::Null => Value::Null,
            WireValue::Integer(i) => Value::Integer(*i),
            WireValue::Real(f) => Value::Real(*f),
            WireValue::Text(s) => Value::Text(s.clone()),
            WireValue::Blob(b) => Value::Blob(decode_blob(b)),
        }
    }

    /// Pipe/format display (NULL→empty, Blob→lowercase hex).
    #[must_use]
    pub(crate) fn format(&self) -> String {
        match self {
            WireValue::Null => String::new(),
            WireValue::Integer(i) => i.to_string(),
            WireValue::Real(f) => f.to_string(),
            WireValue::Text(s) => s.clone(),
            WireValue::Blob(b) => crate::util::hex_string(&decode_blob(b)),
        }
    }
}

#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub(crate) struct QueryRequest {
    pub store: String, // "core" or "logs" (physical store names)
    pub sql: String,
    #[serde(default)]
    pub params: Vec<WireValue>,
}

#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub(crate) struct QueryResponse {
    pub columns: Vec<String>,
    pub rows: Vec<Vec<WireValue>>,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub error: Option<String>,
    #[serde(default)]
    pub truncated: bool,
}

async fn write_frame<S>(stream: &mut S, payload: &[u8]) -> std::io::Result<()>
where
    S: AsyncWrite + Unpin,
{
    let len = u32::try_from(payload.len())
        .map_err(|_| std::io::Error::other("payload too large for IPC frame"))?;
    stream.write_all(&len.to_le_bytes()).await?;
    stream.write_all(payload).await?;
    stream.flush().await
}

async fn read_frame<S>(stream: &mut S) -> std::io::Result<Vec<u8>>
where
    S: AsyncRead + Unpin,
{
    let mut len_buf = [0u8; FRAME_LEN];
    stream.read_exact(&mut len_buf).await?;
    let len = u32::from_le_bytes(len_buf) as usize;
    if len > MAX_FRAME_LEN {
        return Err(std::io::Error::other(format!(
            "IPC frame too large: {len} bytes (max {MAX_FRAME_LEN})"
        )));
    }
    let mut payload = vec![0u8; len];
    stream.read_exact(&mut payload).await?;
    Ok(payload)
}

fn response_from_rows(rows: ReadonlyRows) -> QueryResponse {
    QueryResponse {
        columns: rows.columns,
        rows: rows
            .rows
            .into_iter()
            .map(|row| row.iter().map(WireValue::from_turso).collect())
            .collect(),
        error: None,
        truncated: rows.truncated,
    }
}

fn response_error(msg: &str) -> QueryResponse {
    QueryResponse {
        columns: Vec::new(),
        rows: Vec::new(),
        error: Some(msg.to_string()),
        truncated: false,
    }
}

/// Serialize a `QueryResponse` into an IPC frame and write it. (serde_json
/// emits `null` — not an error — for non-finite floats, which `de_real` maps
/// back to NaN, so serialization never fails for the response types used.)
async fn write_response<S>(stream: &mut S, resp: &QueryResponse) -> std::io::Result<()>
where
    S: AsyncWrite + Unpin,
{
    let payload = serde_json::to_vec(resp).map_err(std::io::Error::other)?;
    write_frame(stream, &payload).await
}

async fn handle_connection(
    mut stream: LocalSocketStream,
    log_store: std::sync::Arc<crate::logs::LogStore>,
) -> anyhow::Result<()> {
    let payload = read_frame(&mut stream).await?;
    let req: QueryRequest = match serde_json::from_slice(&payload) {
        Ok(req) => req,
        Err(e) => {
            // A malformed request must get an actionable error, not an EOF/timeout.
            let resp = response_error(&format!("malformed IPC request: {e}"));
            write_response(&mut stream, &resp).await?;
            return Ok(());
        }
    };

    let conn = match req.store.as_str() {
        "logs" => Some(log_store.conn.clone()),
        "core" => crate::db::DOMAIN_CONN.get().cloned(),
        other => {
            let resp = response_error(&format!(
                "unknown store '{other}' (expected \"core\" or \"logs\")"
            ));
            write_response(&mut stream, &resp).await?;
            return Ok(());
        }
    };

    let resp = match conn {
        Some(conn) => {
            let params: Vec<Value> = req.params.iter().map(WireValue::to_turso).collect();
            match conn.query_readonly(&req.sql, params, IPC_ROW_LIMIT).await {
                Ok(rows) => response_from_rows(rows),
                Err(e) => response_error(&format!("{e}")),
            }
        }
        None => response_error("the requested store is not initialized"),
    };

    write_response(&mut stream, &resp).await?;
    Ok(())
}

/// Spawn the instance-side IPC listener. Call after `DOMAIN_CONN` / `LOG_STORE`
/// are set. Exits (dropping the listener) on shutdown.
pub async fn run_ipc_listener(
    storage_root: &std::path::Path,
    log_store: std::sync::Arc<crate::logs::LogStore>,
) {
    let socket = socket_path(storage_root);

    let name: Name<'static> = match socket.as_path().to_fs_name::<GenericFilePath>() {
        Ok(name) => name.into_owned(),
        Err(e) => {
            warn!(
                error = %e,
                socket = %socket.display(),
                "ipc: cannot build socket name; debug IPC disabled"
            );
            return;
        }
    };

    let mut opts = ListenerOptions::new();
    opts = opts.name(name).try_overwrite(true);

    let listener = match opts.create_tokio() {
        Ok(l) => l,
        Err(e) if e.kind() == std::io::ErrorKind::AddrInUse => {
            warn!(
                error = %e,
                socket = %socket.display(),
                "ipc: socket already in use (another instance?); debug IPC disabled"
            );
            return;
        }
        Err(e) => {
            warn!(
                error = %e,
                socket = %socket.display(),
                "ipc: cannot bind debug socket; debug IPC disabled"
            );
            return;
        }
    };
    // Restrictive perms: only the instance's user can connect. `ListenerOptions::mode`
    // is NOT used — it calls `fchmod` on the socket fd, which returns
    // EINVAL→Unsupported on macOS UDS; `chmod` via the path works everywhere.
    #[cfg(unix)]
    {
        use std::os::unix::fs::PermissionsExt;
        if let Err(e) = std::fs::set_permissions(&socket, std::fs::Permissions::from_mode(0o600)) {
            warn!(
                error = %e,
                socket = %socket.display(),
                "ipc: failed to restrict debug socket perms"
            );
        }
    }
    debug!(socket = %socket.display(), "ipc: debug query listener running");
    let shutdown = crate::shutdown::shutdown_token();

    loop {
        tokio::select! {
            () = shutdown.cancelled() => {
                debug!("ipc: shutdown — closing debug query listener");
                break;
            }
            accepted = listener.accept() => {
                match accepted {
                    Ok(stream) => {
                        let store = log_store.clone();
                        tokio::spawn(async move {
                            if let Err(e) = handle_connection(stream, store).await {
                                debug!(error = %e, "ipc: connection handling error");
                            }
                        });
                    }
                    Err(e) => {
                        warn!(error = %e, "ipc: accept error");
                    }
                }
            }
        }
    }
}

/// Client helper: connect to the instance's IPC socket and run a read-only query.
pub(crate) async fn ipc_query(
    storage_root: &Path,
    req: QueryRequest,
) -> anyhow::Result<QueryResponse> {
    let socket = socket_path(storage_root);
    let name = socket
        .as_path()
        .to_fs_name::<GenericFilePath>()
        .map_err(|e| anyhow::anyhow!("instance IPC endpoint not reachable: {e}"))?;

    let result = tokio::time::timeout(Duration::from_secs(5), async {
        let mut stream = LocalSocketStream::connect(name).await?;
        let payload = serde_json::to_vec(&req)?;
        write_frame(&mut stream, &payload).await?;
        let payload = read_frame(&mut stream).await?;
        let resp: QueryResponse = serde_json::from_slice(&payload)?;
        Ok::<QueryResponse, anyhow::Error>(resp)
    })
    .await
    .map_err(|_| anyhow::anyhow!("instance IPC endpoint timed out"))??;

    Ok(result)
}

/// Backoff delay for retry `attempt` (0-based), clamped to the schedule.
fn retry_backoff(attempt: usize) -> Duration {
    Duration::from_millis(IPC_RETRY_BACKOFF_MS[attempt.min(IPC_RETRY_BACKOFF_MS.len() - 1)])
}

/// Refusal wording for the client side when an instance holds the storage
/// location but its debug channel cannot be reached. Callers only get here after
/// the probe reported the location held, so a bare OS error ("connection
/// refused", "no such file") would read as "the store is gone" instead of "an
/// instance is running and its store was deliberately left alone".
fn channel_unreachable(socket: &Path) -> String {
    format!(
        "a mahbot instance is running against this storage location and holds its stores, but its \
         debug channel at {} could not be reached — the live store was left untouched",
        socket.display()
    )
}

/// Bounded-retry async client used by `mahbot debug` when an instance holds the
/// storage location but the IPC socket is not yet bound (it is between
/// lock-acquire and listener-bind, e.g. boot or self-update hand-off). The
/// caller must only use this when [`crate::util::lock::instance_lock_state_settled`]
/// reports the location [`Held`](crate::util::lock::InstanceLockState::Held) —
/// otherwise the retry would mask a genuinely-down instance.
pub(crate) async fn ipc_query_with_wait(
    storage_root: &Path,
    req: &QueryRequest,
) -> anyhow::Result<QueryResponse> {
    let socket = socket_path(storage_root);
    let deadline = std::time::Instant::now() + ipc_bound_timeout();
    let mut attempt = 0usize;
    loop {
        match ipc_query(storage_root, req.clone()).await {
            Ok(resp) => return Ok(resp),
            Err(_) if std::time::Instant::now() < deadline => {
                tokio::time::sleep(retry_backoff(attempt)).await;
                attempt += 1;
            }
            Err(e) => return Err(e.context(channel_unreachable(&socket))),
        }
    }
}

/// Synchronous bounded-retry client used by `bench-openrouter`'s synchronous
/// config-resolution path (which cannot await the async [`ipc_query_with_wait`];
/// the call blocks the current thread). Like the async client, the caller must
/// only use it after the instance-lock probe reported the location held.
pub(crate) fn ipc_query_sync(
    storage_root: &Path,
    req: &QueryRequest,
) -> std::io::Result<QueryResponse> {
    let socket = socket_path(storage_root);
    let deadline = std::time::Instant::now() + ipc_bound_timeout();
    let mut attempt = 0usize;
    loop {
        match ipc_query_sync_once(storage_root, req) {
            Ok(resp) => return Ok(resp),
            Err(_) if std::time::Instant::now() < deadline => {
                std::thread::sleep(retry_backoff(attempt));
                attempt += 1;
            }
            // The bare OS text would read as "the store is gone"; our sentence
            // says the instance is running and its store was left alone. The
            // failure kind is not carried through — no caller distinguishes it.
            Err(e) => {
                return Err(std::io::Error::other(format!(
                    "{}: {e}",
                    channel_unreachable(&socket)
                )));
            }
        }
    }
}

fn ipc_query_sync_once(storage_root: &Path, req: &QueryRequest) -> std::io::Result<QueryResponse> {
    let socket = socket_path(storage_root);
    let name = socket
        .as_path()
        .to_fs_name::<GenericFilePath>()
        .map_err(|e| std::io::Error::other(format!("instance IPC endpoint not reachable: {e}")))?;
    // `Stream::connect` is an associated function on the sync `Stream` trait,
    // not the enum — call it fully-qualified so `Self` resolves to the enum.
    let mut stream =
        <interprocess::local_socket::Stream as interprocess::local_socket::traits::Stream>::connect(
            name,
        )?;
    let payload = serde_json::to_vec(req).map_err(std::io::Error::other)?;
    write_frame_sync(&mut stream, &payload)?;
    let payload = read_frame_sync(&mut stream)?;
    let resp: QueryResponse = serde_json::from_slice(&payload).map_err(std::io::Error::other)?;
    Ok(resp)
}

fn write_frame_sync(stream: &mut impl std::io::Write, payload: &[u8]) -> std::io::Result<()> {
    let len = u32::try_from(payload.len())
        .map_err(|_| std::io::Error::other("payload too large for IPC frame"))?;
    stream.write_all(&len.to_le_bytes())?;
    stream.write_all(payload)?;
    stream.flush()
}

fn read_frame_sync(stream: &mut impl std::io::Read) -> std::io::Result<Vec<u8>> {
    let mut len_buf = [0u8; FRAME_LEN];
    stream.read_exact(&mut len_buf)?;
    let len = u32::from_le_bytes(len_buf) as usize;
    if len > MAX_FRAME_LEN {
        return Err(std::io::Error::other(format!(
            "IPC frame too large: {len} bytes (max {MAX_FRAME_LEN})"
        )));
    }
    let mut payload = vec![0u8; len];
    stream.read_exact(&mut payload)?;
    Ok(payload)
}

#[cfg(test)]
mod tests {
    use super::*;

    /// The IPC `WireValue` round-trips through `from_turso`/`to_turso` and
    /// renders NULL as empty, integers as decimals, and blobs as lowercase hex.
    #[test]
    fn wire_value_round_trips() {
        use crate::db::Value;
        for (value, wire) in [
            (Value::Integer(42), WireValue::Integer(42)),
            (Value::Real(1.5), WireValue::Real(1.5)),
            (
                Value::Text("hi".to_string()),
                WireValue::Text("hi".to_string()),
            ),
            (Value::Null, WireValue::Null),
        ] {
            assert_eq!(WireValue::from_turso(&value), wire, "from_turso");
            assert_eq!(wire.to_turso(), value, "to_turso");
        }
        // Blob round-trips through the base64 wire form.
        let blob = vec![0x00u8, 0xDE, 0xAD];
        let wire = WireValue::from_turso(&Value::Blob(blob.clone()));
        assert_eq!(wire.to_turso(), Value::Blob(blob));
        assert_eq!(wire.format(), "00dead");
        assert_eq!(WireValue::Null.format(), "");
    }

    /// serde_json emits `null` for non-finite floats; `de_real` must map it back
    /// to `NaN` so a REAL NaN cell round-trips instead of failing deserialization.
    #[test]
    fn wire_value_round_trips_non_finite_real() {
        let payload = serde_json::to_vec(&WireValue::Real(f64::NAN)).unwrap();
        match serde_json::from_slice::<WireValue>(&payload).unwrap() {
            WireValue::Real(f) => {
                assert!(f.is_nan(), "non-finite REAL must round-trip as NaN");
            }
            other => panic!("expected Real, got {other:?}"),
        }
    }

    /// End-to-end: the instance-side listener serves a read-only query over the
    /// local socket, returning column names + rows (the IPC path `mahbot debug`
    /// uses when an instance holds the location). `serial(ipc_bound)`: this
    /// relies on the bounded retry (the listener task may bind after the first
    /// attempt) and must not overlap the tests that shrink that bound through
    /// `MAHBOT_IPC_BOUND_TIMEOUT_SECS`.
    #[serial_test::serial(ipc_bound)]
    #[tokio::test]
    async fn ipc_query_serves_readonly_queries_end_to_end() {
        let (store, dir) = crate::open_test_store!(crate::logs::LogStore, "log");
        let store_arc = std::sync::Arc::new(store);
        // Insert a row so COUNT is non-zero (proves the query reads real data).
        store_arc
            .conn
            .execute_batch(
                "INSERT INTO logs (timestamp, level, target, message) \
                 VALUES ('2026-01-01T00:00:00Z', 'INFO', 'test', 'hello');",
            )
            .await
            .unwrap();

        let root = dir.path().to_path_buf();
        let listener_root = root.clone();
        let listener = tokio::spawn(async move {
            crate::db::ipc::run_ipc_listener(&listener_root, store_arc).await;
        });

        let req = QueryRequest {
            store: "logs".to_string(),
            sql: "SELECT COUNT(*) FROM logs".to_string(),
            params: Vec::new(),
        };
        let resp = crate::db::ipc::ipc_query_with_wait(&root, &req)
            .await
            .expect("IPC query must reach the listener");
        assert!(resp.error.is_none(), "no error expected: {:?}", resp.error);
        assert_eq!(resp.columns, vec!["COUNT(*)"]);
        assert_eq!(resp.rows, vec![vec![WireValue::Integer(1)]]);
        assert!(!resp.truncated);

        // Engine write enforcement: a write statement is rejected with
        // query_only=ON (defense-in-depth on top of the CLI blocklist).
        let write_req = QueryRequest {
            store: "logs".to_string(),
            sql: "CREATE TABLE nope (id INTEGER)".to_string(),
            params: Vec::new(),
        };
        let write_resp = crate::db::ipc::ipc_query_with_wait(&root, &write_req)
            .await
            .expect("IPC write attempt must reach the listener");
        assert!(
            write_resp.error.is_some(),
            "a write statement must be rejected by query_only=ON"
        );

        // The shared connection's query_only must be reset after the query so
        // the instance's writes are never left disabled.
        let reset_check = crate::db::ipc::ipc_query_with_wait(
            &root,
            &QueryRequest {
                store: "logs".to_string(),
                sql: "SELECT 1".to_string(),
                params: Vec::new(),
            },
        )
        .await
        .expect("subsequent query must succeed after reset");
        assert!(reset_check.error.is_none());

        listener.abort();
        let _ = listener.await;
    }
}