Compare commits
9 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| de026b8f8a | |||
| 97f6c33a45 | |||
| 7d4a4339c2 | |||
| a55079afbd | |||
| e18ad4723b | |||
| 6de8ac8972 | |||
| a052975420 | |||
| 3743789e8d | |||
| ee5f8b932b |
+166
@@ -0,0 +1,166 @@
|
||||
# CB-137 / fleetd issue #137 — report
|
||||
|
||||
## Real root cause (not the hypothesis in the ticket)
|
||||
|
||||
I read `MessageService.java` and `Rendezvous.java` before changing anything. The mechanism is real,
|
||||
but the exact place it happens is `MessageService.answer()`, not "the reply goes to the inbox on
|
||||
purpose" in general.
|
||||
|
||||
1. A lead delegates with `fleet_send{wait:false}` → `sendAsync()` creates a `Task` and runs `send()`
|
||||
on a background virtual thread with a 30-minute internal budget (`ASYNC_TIMEOUT_MS`).
|
||||
2. The worker calls `fleet_ask`. That resolves the open rendezvous waiter with `Kind.QUESTION`, so
|
||||
`send()` returns immediately and the `Task` is left open (its `future` stays unresolved — see the
|
||||
comment in `sendAsync`'s lambda: "Keep the accepted owner until answer() finishes it").
|
||||
3. The lead answers with `fleet_send{turnId, content}`. This calls `FleetMcp.answer()` →
|
||||
`MessageService.answer(turnId, content, timeout)`. The `timeout` here is **not** the generous
|
||||
30-minute async budget — it is the MCP tool's own bounded wait: `DEFAULT_TIMEOUT_MS = 25_000`,
|
||||
clamped to at most `MAX_TIMEOUT_MS = 120_000` (`FleetMcp.java:71-72,495,512`). This is the same
|
||||
~60–120s window every blocking `fleet_send` call is capped at (documented elsewhere as "the
|
||||
caller's own MCP client call timeout").
|
||||
4. `answer()` opens a **fresh** rendezvous waiter for the worker session and blocks on it for at most
|
||||
that window. If the worker's resumed turn takes longer than that to actually finish (very
|
||||
plausible — the resumed turn can mean more edits, a build, a commit, a push, opening a PR), the
|
||||
wait times out. On timeout, `answer()`'s `finally` block unconditionally calls
|
||||
`rendezvous.close(workerSession, reply)`, **removing the waiter from the map**, and returns
|
||||
`Outcome.TIMED_OUT_WORKING` to the lead.
|
||||
5. The worker keeps working, unaware anything happened, and eventually calls `fleet_reply`. That
|
||||
reaches `MessageService.reply(session, content)`, which tries `rendezvous.resolve(session,
|
||||
content)` — but the waiter was already closed in step 4, so `resolve` returns `false`. `reply()`
|
||||
then falls back to `inbox.publish(...)` and marks `strandedReplies.put(session, true)`
|
||||
(CB-640 bookkeeping) — the reply is safely held, but **the async `Task`'s `future` is never
|
||||
completed**.
|
||||
6. `fleet_poll{ticket}` keeps returning `PENDING` forever (the `Task` never resolves) — until the
|
||||
lead eventually calls `fleet_stop`. That fires `sessions.onRelease` → `messages.abandon(target,
|
||||
reason)` (`Fleetd.java:481-497`), where `reason` is built with the exact text from the bug report
|
||||
("the worker session was released before it replied; worktree=... branch=... snapshot=...",
|
||||
`Fleetd.java:484-487`). `abandon()`'s loop finds the still-open `Task` (`question == null`, future
|
||||
not done) and completes it as `WORKER_FAILED` with that misleading reason — even though the
|
||||
worker's real reply is sitting, intact, in the inbox the whole time.
|
||||
|
||||
So: the reported behaviour is correct, and the specific trigger is `answer()`'s own bounded wait
|
||||
being shorter than the worker's real resumed-turn time — not anything to do with the ~55s
|
||||
`fleet_ask` window itself (that part, issue #61, is untouched).
|
||||
|
||||
## Fix
|
||||
|
||||
Two changes in `fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java`, both scoped to the
|
||||
ticket/reply routing and the terminal-state text — `fleet_ask`'s own window and mechanics are
|
||||
untouched.
|
||||
|
||||
**1. `reply()` — priority 1 (the ticket resolves with the real reply).**
|
||||
Before falling back to the inbox, `reply()` now looks for an async `Task` that is specifically in the
|
||||
"already answered but not yet resolved" state (`question == null`, `turnId != null` — set once
|
||||
`answer()` has cleared the question but before anything completed the future, `!future.isDone()`).
|
||||
If one exists for this `target`, the worker's reply completes that `Task`'s future directly as
|
||||
`Outcome.REPLIED` with the real content, and the reply never touches the inbox at all. A task that
|
||||
was never asked has `turnId == null` and can never match, so ordinary (no-`fleet_ask`) delegations
|
||||
are unaffected — they already resolve through the pre-existing rendezvous fast path.
|
||||
|
||||
I chose this over leaving `answer()`'s own timeout behaviour untouched and instead keeping its
|
||||
rendezvous waiter open in the background: that alternative works but reopens the "at most one
|
||||
waiter per session" invariant (`Rendezvous.open` throws on a double-open) to a new class of races
|
||||
with a fresh send arriving mid-window. The `send()` path already guards against sending into an
|
||||
answered-but-still-resolving worker via `hasAsyncQuestion(target)` (checks `asyncTasksByTurn`,
|
||||
which still holds the task until it resolves), so routing through `reply()` gets the same protection
|
||||
without touching `answer()`'s waiter lifecycle at all — the smaller, safer diff.
|
||||
|
||||
**2. `abandon()` — priority 3 (required independently, "even if you fix (1)").**
|
||||
Before marking any of a released target's still-open tasks `WORKER_FAILED`, `abandon()` now checks
|
||||
`hasStrandedReply(target)` (the existing CB-640 fact — true whenever the *last* `reply()` for this
|
||||
target fell through to the inbox). If true, it drains the inbox (`recoverStrandedReply`) and — if it
|
||||
actually finds a message — completes the task as `REPLIED` with that real content instead of writing
|
||||
the failure. This is deliberately a **separate** check from fix 1: fix 1 already prevents the
|
||||
inbox-stranding from happening in the exact scenario this ticket describes, so by the time
|
||||
`abandon()` runs the task is normally already resolved and `abandon()`'s `complete()` call is a
|
||||
harmless no-op. This second check exists so that if some *other* future path ever strands a reply
|
||||
in the inbox without resolving its ticket, `abandon()` still refuses to report a false failure —
|
||||
"if a reply reached any sink for that turn, the terminal state is done," per the ticket. I verified
|
||||
both are required by disabling each independently and confirming the two new tests fail (see below).
|
||||
|
||||
**Priority 4 (the snapshot/worktree hint).** Handled as a consequence of both fixes rather than a
|
||||
separate branch: once a task resolves as `REPLIED` (via either fix), `abandon()` never calls
|
||||
`new Reply(Outcome.WORKER_FAILED, reason)` for that task at all, so the "the worker session was
|
||||
released before it replied; worktree=... branch=... snapshot=..." text is never constructed or
|
||||
attached to that ticket's outcome. It still appears, correctly, for a task that never got a reply
|
||||
(the existing `abandonFailsEveryPendingAsyncTicketForTheReleasedTarget` /
|
||||
`anAbandonedAsyncTaskPollsAsFailedNotPending` tests still pass unchanged).
|
||||
|
||||
**Priority 2** was not needed — fix 1 makes `fleet_poll{ticket}` return the actual reply (the
|
||||
higher-priority option), so I did not fall back to "the ticket merely resolves as done with no
|
||||
content."
|
||||
|
||||
## Tests — driven through the real delegation path, not the reply sink directly
|
||||
|
||||
Both new tests in `fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java` go through
|
||||
`sendAsync` → `injectDelivery` → `ask` → `answer` (with a short timeout, so it genuinely times out,
|
||||
mirroring the ~25–120s real MCP-call bound vs. a longer resumed turn) → `reply` → `poll`/`abandon`.
|
||||
No test constructs a `Reply` and hands it to a sink directly.
|
||||
|
||||
- `aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket` — asserts `fleet_poll{ticket}` (via
|
||||
`messages.poll`) reaches `Phase.DONE` with the worker's actual reply text and
|
||||
`replySource() == "reply"`, and that `hasStrandedReply(T)` stays `false` (proves the reply never
|
||||
touched the inbox at all — fix 1 caught it).
|
||||
- `fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket` — same setup, then calls `abandon(T, "the
|
||||
worker session was released before it replied")` (what `fleet_stop` triggers) and asserts it
|
||||
returns `false` (no failure recorded) and the ticket still polls `DONE` with the real reply.
|
||||
|
||||
**Proof both fail without the change.** I temporarily short-circuited both new private methods
|
||||
(`askAnsweredAsyncTask` → always `null`, `recoverStrandedReply` → always `null`) — i.e. disabled
|
||||
both fixes — and ran just these two tests:
|
||||
|
||||
```
|
||||
[ERROR] Tests run: 2, Failures: 2, Errors: 0, Skipped: 0
|
||||
dev.ltms.fleet.msg.MessageServiceTest.aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket
|
||||
org.opentest4j.AssertionFailedError: expected: <DONE> but was: <PENDING>
|
||||
dev.ltms.fleet.msg.MessageServiceTest.fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket
|
||||
org.opentest4j.AssertionFailedError: a reply already arrived, so nothing here is a genuine failure
|
||||
==> expected: <false> but was: <true>
|
||||
```
|
||||
|
||||
This is the exact bug: the ticket stays `PENDING` forever, and `abandon()` reports `true` (a
|
||||
failure) even though a reply had already arrived. I then restored both fixes (verified with
|
||||
`grep -n "TEMP #137-proof"` finding nothing) and re-ran — both pass.
|
||||
|
||||
## Build
|
||||
|
||||
Ran from `fleetd/`, unpiped, full output read (not `| tail`):
|
||||
|
||||
```
|
||||
mvn clean install
|
||||
...
|
||||
[INFO] Tests run: 1039, Failures: 0, Errors: 0, Skipped: 0
|
||||
[INFO] BUILD SUCCESS
|
||||
[INFO] Total time: 36.315 s
|
||||
```
|
||||
|
||||
Main was at 1037 tests; this branch adds the 2 new tests above → 1039, all green, `exit=0`.
|
||||
|
||||
## What I could NOT check
|
||||
|
||||
- No IDE tooling is mounted for me (worker), so no `ide_diagnostics`/IntelliJ inspection pass — only
|
||||
`mvn clean install` (compiler + full test suite), as the worker procedure allows.
|
||||
- I cannot restart the daemon or dogfood this live — I have no forge/daemon control. This is
|
||||
unverified against a real herdr pane, a real MCP client's ~60s call cap, or a real worker session;
|
||||
everything above is verified only through the JUnit fixture's simulated timing
|
||||
(`FakeHerdr`/`injector.onStatus`/direct `messages.answer(...,150)` calls), not a live fleet.
|
||||
A primary should still consider a short live dogfood (an async delegation that asks, gets answered,
|
||||
and takes longer than ~2 minutes to reply) before calling this closed.
|
||||
- I did not touch, and did not re-verify, the `fleet_ask` ~55s window itself (issue #61) — out of
|
||||
scope per the brief.
|
||||
|
||||
## Scope note (not investigated further)
|
||||
|
||||
`answer()`'s nested/double-`fleet_ask` case (the worker asks a second question before ever
|
||||
replying to the first answer) has some pre-existing behaviour around which `turnId` a `QUESTION`
|
||||
resolution gets attributed to that I did not fully untangle — it predates this change, my fix does
|
||||
not touch it, and it is unrelated to the reported defect. Flagging only; not investigated further.
|
||||
|
||||
## Handoff
|
||||
|
||||
- Branch: `worker/cb-137-ask-ticket-e7760c-2`
|
||||
- Worktree root: `/Users/dai.ha/LTMS/.bridged-worktrees/734324-2`
|
||||
- Files changed:
|
||||
- `fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java`
|
||||
- `fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java`
|
||||
- `REPORT-cb137.md` (this file)
|
||||
- Build: `Tests run: 1039, Failures: 0, Errors: 0, Skipped: 0` / `BUILD SUCCESS` (verbatim above)
|
||||
@@ -28,6 +28,7 @@
|
||||
<testcontainers.version>1.20.4</testcontainers.version>
|
||||
<commons-compress.version>1.27.1</commons-compress.version>
|
||||
<commons-lang3.version>3.18.0</commons-lang3.version>
|
||||
<sqlite-jdbc.version>3.53.4.0</sqlite-jdbc.version>
|
||||
</properties>
|
||||
|
||||
<!--
|
||||
@@ -44,6 +45,12 @@
|
||||
3.0-rc5; bumping Jackson 3 to the patched 3.2.x breaks the SDK (annotation mismatch).
|
||||
Only the loopback /mcp endpoint parses this JSON, from trusted local Claude clients.
|
||||
The 11.0.23 -> 11.0.25 bump did clear jetty CVE-2024-8184 (5.9) and CVE-2024-6763.
|
||||
|
||||
fleetd #206: org.xerial:sqlite-jdbc 3.53.4.0 (added for OpenCodeSessionDiscovery) — the
|
||||
only known advisory against this artifact is CVE-2023-32697 (RCE via an attacker-controlled
|
||||
JDBC URL), fixed in 3.41.2.2; 3.53.4.0 is well past that fix and OSV.dev reports no open
|
||||
advisory against it. Checked via the OSV.dev API (no Mend.io/JetBrains IDE MCP mount
|
||||
available from this worktree) on 2026-08-31.
|
||||
-->
|
||||
|
||||
<!-- Force the latest patched Jetty 11.x across all Javalin-pulled Jetty modules (no version
|
||||
@@ -123,6 +130,17 @@
|
||||
<version>${amqp.version}</version>
|
||||
</dependency>
|
||||
|
||||
<!-- fleetd #206: opencode moved its session store from a JSON tree to SQLite
|
||||
(opencode.db). This is the JDBC driver OpenCodeSessionDiscovery uses to read it
|
||||
read-only. Ships bundled native libraries (linux/mac/windows, several archs), so it
|
||||
is a heavier jar than most deps here — see the pom's dependency-security note below
|
||||
for the size/CVE tradeoff actually measured. -->
|
||||
<dependency>
|
||||
<groupId>org.xerial</groupId>
|
||||
<artifactId>sqlite-jdbc</artifactId>
|
||||
<version>${sqlite-jdbc.version}</version>
|
||||
</dependency>
|
||||
|
||||
<!-- Logging -->
|
||||
<dependency>
|
||||
<groupId>org.slf4j</groupId>
|
||||
|
||||
@@ -1,58 +1,87 @@
|
||||
package dev.ltms.fleet.member;
|
||||
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.sqlite.SQLiteConfig;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.stream.Stream;
|
||||
import java.sql.Connection;
|
||||
import java.sql.PreparedStatement;
|
||||
import java.sql.ResultSet;
|
||||
import java.sql.SQLException;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
|
||||
/**
|
||||
* Resolves the opencode session id for a fleetd worker from opencode's on-disk storage — the
|
||||
* only place this adapter touches opencode's private layout, and deliberately the <em>only</em>
|
||||
* class that does.
|
||||
*
|
||||
* <p><strong>Why this is isolated behind one seam.</strong> The layout is version-coupled and not a
|
||||
* stable contract: opencode writes one JSON file per session under
|
||||
* {@code <storageRoot>/session/<projectID>/<ses_*.json>}, and each record carries a
|
||||
* {@code "version"} field (e.g. {@code "1.1.31"}), so the exact directory shape, file naming, and
|
||||
* field names can move between opencode releases. opencode also ships a headless HTTP server that
|
||||
* may supersede file scanning entirely. Everything this adapter knows about that private storage —
|
||||
* its shape, naming, and field names — lives here, so a layout change, or a switch to the HTTP
|
||||
* server, changes exactly one class and nothing in {@link OpenCodeLauncher}.
|
||||
* <p><strong>Why this is isolated behind one seam.</strong> The layout is version-coupled and not
|
||||
* a stable contract: opencode persists its session state in a SQLite database at
|
||||
* {@code <storageRoot>/opencode.db} (a {@code session} table, one row per session, keyed by id and
|
||||
* carrying a {@code directory} column). That schema can move between opencode releases exactly
|
||||
* like the JSON-file layout it replaced did (opencode migrated off a one-JSON-file-per-session
|
||||
* tree under {@code <storageRoot>/storage/session/<projectID>/ses_*.json} in January 2026 — that
|
||||
* tree is now a frozen migration artefact nothing writes, which is why this class no longer reads
|
||||
* it). opencode also ships a headless HTTP server that may supersede both of these entirely.
|
||||
* Everything this adapter knows about that private storage — its shape and column names — lives
|
||||
* here, so a layout change, or a switch to the HTTP server, changes exactly one class and nothing
|
||||
* in {@link OpenCodeLauncher}.
|
||||
*
|
||||
* <p>The determinism that makes this useful is structural, not a guess: every fleetd worker runs
|
||||
* in its own unique git worktree, so the record's {@code directory} (its project root) equals the
|
||||
* in its own unique git worktree, so the row's {@code directory} (its project root) equals the
|
||||
* worker's cwd identifies <em>its</em> session unambiguously. We match on {@code directory} rather
|
||||
* than diffing {@code opencode session list} before/after — that races under concurrent spawns, and
|
||||
* the CLI listing does not even show the directory.
|
||||
*
|
||||
* <p>All reads are best-effort and never throw: a missing or unreadable storage root, a record that
|
||||
* fails to parse, or a directory with no record yet all yield {@code null}, and the caller (the
|
||||
* session handle) treats that as "identity not resolved yet" and retries later.
|
||||
* <p>All reads are best-effort and never throw: a missing or unreadable database, a query that
|
||||
* fails, or a directory with no row yet all yield {@code null}, and the caller (the session
|
||||
* handle) treats that as "identity not resolved yet" and retries later. The database is opened
|
||||
* read-only and never written to: opencode itself may be running and writing it concurrently (WAL
|
||||
* mode), and this class must never disturb that.
|
||||
*/
|
||||
final class OpenCodeSessionDiscovery {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(OpenCodeSessionDiscovery.class);
|
||||
|
||||
private final Path storageRoot; // e.g. ~/.local/share/opencode (injectable for tests)
|
||||
private final ObjectMapper json;
|
||||
private final Path databasePath;
|
||||
private final AtomicBoolean warnedMissingDatabase = new AtomicBoolean(false);
|
||||
|
||||
OpenCodeSessionDiscovery(Path storageRoot) {
|
||||
this.storageRoot = storageRoot;
|
||||
this.json = new ObjectMapper();
|
||||
this.databasePath = storageRoot.resolve("opencode.db");
|
||||
}
|
||||
|
||||
/**
|
||||
* The opencode session id whose record references {@code directory} (the worker's cwd), or
|
||||
* {@code null} when no record matches yet. When several records share the directory — e.g.
|
||||
* repeated spawns into the same worktree — the <em>most recently modified</em> one wins: it is
|
||||
* A connection to {@link #databasePath} opened with SQLite's {@code SQLITE_OPEN_READONLY}
|
||||
* flag: it never creates the file, never writes, and never touches WAL or journal mode.
|
||||
* opencode may be running and writing this database concurrently, and this class must never
|
||||
* disturb it.
|
||||
*
|
||||
* <p>Package-private so a test can hold the connection and prove it refuses a write. That is
|
||||
* the only way to pin this property: making the file unwritable does <em>not</em> work,
|
||||
* because SQLite silently downgrades a read-write open of an unwritable file to read-only, so
|
||||
* such a test passes whether or not the flag is set.
|
||||
*/
|
||||
Connection openReadOnly() throws SQLException {
|
||||
SQLiteConfig config = new SQLiteConfig();
|
||||
config.setReadOnly(true);
|
||||
return config.createConnection("jdbc:sqlite:" + databasePath);
|
||||
}
|
||||
|
||||
/**
|
||||
* The opencode session id whose row references {@code directory} (the worker's cwd), or
|
||||
* {@code null} when no row matches yet. When several rows share the directory — e.g. repeated
|
||||
* spawns into the same worktree — the row with the highest {@code time_updated} wins: it is
|
||||
* the session the pane most likely corresponds to.
|
||||
*
|
||||
* <p>Never throws: a missing {@code storageRoot}, an unreadable/malformed record, or a
|
||||
* directory that has not been persisted yet all resolve to {@code null} rather than failing a
|
||||
* spawn. A fleetd worker's session record is written lazily (when the session is first
|
||||
* persisted), so {@code null} here is the normal answer right after the pane is ready, and the
|
||||
* caller retries later.
|
||||
* <p>Never throws: a missing {@code opencode.db}, a locked/unreadable database, a query
|
||||
* failure, or a directory that has not been persisted yet all resolve to {@code null} rather
|
||||
* than failing a spawn. A fleetd worker's session row is written lazily (when the session is
|
||||
* first persisted), so {@code null} here is the normal answer right after the pane is ready,
|
||||
* and the caller retries later.
|
||||
*
|
||||
* @param directory the worker's cwd, as resolved for this spawn
|
||||
* @return the matching session id, or {@code null} if none is known yet
|
||||
@@ -61,62 +90,31 @@ final class OpenCodeSessionDiscovery {
|
||||
if (directory == null || directory.isBlank()) {
|
||||
return null;
|
||||
}
|
||||
Path sessionRoot = storageRoot.resolve("session");
|
||||
if (!Files.isDirectory(sessionRoot)) {
|
||||
if (!Files.isRegularFile(databasePath)) {
|
||||
if (warnedMissingDatabase.compareAndSet(false, true)) {
|
||||
log.warn("opencode session database not found at {} — opencode's on-disk layout "
|
||||
+ "may have moved again; session discovery will keep returning null",
|
||||
databasePath);
|
||||
}
|
||||
return null;
|
||||
}
|
||||
String best = null;
|
||||
long bestMtime = Long.MIN_VALUE;
|
||||
try (Stream<Path> projectDirs = Files.list(sessionRoot)) {
|
||||
for (Path projectDir : projectDirs.filter(Files::isDirectory).toList()) {
|
||||
try (Stream<Path> records = Files.list(projectDir)) {
|
||||
for (Path record : records.toList()) {
|
||||
String id = matchId(record, directory);
|
||||
if (id == null) {
|
||||
continue;
|
||||
}
|
||||
long mtime = lastModifiedEpochMillis(record);
|
||||
if (mtime > bestMtime) {
|
||||
bestMtime = mtime;
|
||||
best = id;
|
||||
}
|
||||
}
|
||||
} catch (IOException ignored) {
|
||||
// one project dir unreadable — skip it; another may still match
|
||||
String sql = "SELECT id FROM session WHERE directory = ? ORDER BY time_updated DESC LIMIT 1";
|
||||
try (Connection connection = openReadOnly();
|
||||
PreparedStatement statement = connection.prepareStatement(sql)) {
|
||||
statement.setString(1, directory);
|
||||
try (ResultSet rows = statement.executeQuery()) {
|
||||
if (rows.next()) {
|
||||
return rows.getString("id");
|
||||
}
|
||||
}
|
||||
} catch (IOException ignored) {
|
||||
// storage root vanished or became unreadable — "no session known yet"
|
||||
} catch (SQLException e) {
|
||||
// Locked, corrupt, or otherwise unreadable — never fatal to a spawn. Not the
|
||||
// "database moved" signal (the file exists), so this stays below WARN.
|
||||
log.debug("opencode session database unreadable at {}: {}", databasePath, e.toString());
|
||||
return null;
|
||||
}
|
||||
return best;
|
||||
}
|
||||
|
||||
/**
|
||||
* The record's session id when it references {@code directory}, else {@code null}. A record
|
||||
* that is not JSON, lacks {@code id}/{@code directory}, or points at a different directory is
|
||||
* simply not our session; a malformed one is skipped, never fatal.
|
||||
*/
|
||||
private String matchId(Path record, String directory) {
|
||||
try {
|
||||
JsonNode node = json.readTree(record.toFile());
|
||||
JsonNode id = node == null ? null : node.get("id");
|
||||
JsonNode dir = node == null ? null : node.get("directory");
|
||||
if (id == null || dir == null || !directory.equals(dir.asText())) {
|
||||
return null;
|
||||
}
|
||||
return id.asText();
|
||||
} catch (IOException e) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/** The record's last-modified epoch ms, or {@code Long.MIN_VALUE} if unreadable (never wins). */
|
||||
private static long lastModifiedEpochMillis(Path record) {
|
||||
try {
|
||||
return Files.getLastModifiedTime(record).toMillis();
|
||||
} catch (IOException e) {
|
||||
return Long.MIN_VALUE;
|
||||
}
|
||||
log.debug("no opencode session row for directory (root={}, directory={})",
|
||||
storageRoot, directory);
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -9,6 +9,7 @@ import dev.ltms.fleet.metrics.Metrics;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.UUID;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
@@ -167,6 +168,12 @@ public final class MessageService {
|
||||
private static final class Task {
|
||||
private final String ticket;
|
||||
private final String target;
|
||||
/**
|
||||
* When this task was created (#137 fix): the tiebreaker for which of several open tasks on
|
||||
* one target gets a recovered reply in {@link #abandon} — the oldest, since it is the one
|
||||
* that has been waiting longest.
|
||||
*/
|
||||
private final long createdNanos;
|
||||
private final CompletableFuture<Reply> future = new CompletableFuture<>();
|
||||
/**
|
||||
* When {@link #future} resolved, or {@code null} while it is still pending — the clock
|
||||
@@ -184,6 +191,7 @@ public final class MessageService {
|
||||
private Task(String ticket, String target, LongSupplier nowNanos) {
|
||||
this.ticket = ticket;
|
||||
this.target = target;
|
||||
this.createdNanos = nowNanos.getAsLong();
|
||||
future.whenComplete((reply, ex) -> completedNanos = nowNanos.getAsLong());
|
||||
}
|
||||
}
|
||||
@@ -375,6 +383,19 @@ public final class MessageService {
|
||||
* {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions
|
||||
* are interactive and must never be queued.
|
||||
*
|
||||
* <p><strong>Ambiguous match also falls to the inbox.</strong> {@link #askAnsweredAsyncTasks}
|
||||
* cannot actually return more than one entry today (see its own javadoc for why — in short,
|
||||
* {@link #hasAsyncQuestion} keeps a target BUSY, so no second task can reach this state, for as
|
||||
* long as an earlier one's {@code turnId} is still stamped). That is an emergent guarantee from
|
||||
* two other facts, not one this method enforces, so this branch stays in as defence in depth
|
||||
* rather than being removed as dead code: if it ever weakens, returning whichever candidate a
|
||||
* {@code ConcurrentHashMap} iteration reaches first would let a genuine reply complete the
|
||||
* <em>wrong</em> ticket — silently handing the lead something that reads like a correct answer to
|
||||
* a delegation the worker never touched, which is worse than a failure because the lead acts on
|
||||
* it. When more than one candidate exists, guessing is not safe: fall back to the inbox exactly
|
||||
* as the zero-candidate case does, and let {@link #abandon} apply the eventual recovery
|
||||
* deterministically instead.
|
||||
*
|
||||
* @return always {@code true} — the reply resolved a live send, completed a parked ticket, or
|
||||
* was queued
|
||||
*/
|
||||
@@ -392,13 +413,21 @@ public final class MessageService {
|
||||
// FAILED with a misleading "session released before it replied" reason, even though the reply
|
||||
// had, in fact, arrived. Completing the matching ticket directly here means fleet_poll{ticket}
|
||||
// sees the real reply instead.
|
||||
Task orphan = askAnsweredAsyncTask(session);
|
||||
if (orphan != null && orphan.future.complete(new Reply(Outcome.REPLIED, content))) {
|
||||
if (orphan.turnId != null) {
|
||||
asyncTasksByTurn.remove(orphan.turnId, orphan);
|
||||
List<Task> candidates = askAnsweredAsyncTasks(session);
|
||||
if (candidates.size() == 1) {
|
||||
Task orphan = candidates.get(0);
|
||||
if (orphan.future.complete(new Reply(Outcome.REPLIED, content))) {
|
||||
if (orphan.turnId != null) {
|
||||
asyncTasksByTurn.remove(orphan.turnId, orphan);
|
||||
}
|
||||
count(FleetMetrics.REPLIES, "path", "async-recovered");
|
||||
return true; // the ticket itself took it — no inbox stranding at all
|
||||
}
|
||||
count(FleetMetrics.REPLIES, "path", "async-recovered");
|
||||
return true; // the ticket itself took it — no inbox stranding at all
|
||||
} else if (candidates.size() > 1) {
|
||||
List<String> tickets = candidates.stream().map(t -> t.ticket).toList();
|
||||
log.warn("reply from {} matches {} open async tickets {} — cannot tell which one it "
|
||||
+ "answers, queuing to the inbox instead of guessing", session, candidates.size(),
|
||||
tickets);
|
||||
}
|
||||
inbox.publish(session, UUID.randomUUID().toString(), content);
|
||||
// CB-640: record the stranding itself (not just the reply text) so fleet health can see a
|
||||
@@ -414,21 +443,38 @@ public final class MessageService {
|
||||
}
|
||||
|
||||
/**
|
||||
* The still-open async task on {@code target} whose {@code fleet_ask} was already answered — its
|
||||
* {@link Task#turnId} is stamped but its {@link Task#question} was cleared by {@link #answer} —
|
||||
* yet whose future is not resolved yet (#137). {@code null} if no such task exists, including the
|
||||
* Every still-open async task on {@code target} whose {@code fleet_ask} was already answered —
|
||||
* its {@link Task#turnId} is stamped but its {@link Task#question} was cleared by {@link #answer}
|
||||
* — yet whose future is not resolved yet (#137). Empty if no such task exists, including the
|
||||
* common case where {@code target}'s worker never used {@code fleet_ask} at all (a task that was
|
||||
* never asked has {@code turnId == null}, so it can never match here and only ever completes
|
||||
* through the ordinary rendezvous fast path in {@link #reply}).
|
||||
*
|
||||
* <p><strong>Returns at most one entry today — verified, not assumed.</strong> {@link #send}
|
||||
* refuses to open a waiter on {@code target} while {@link #hasAsyncQuestion} is true, and that
|
||||
* check matches ANY task whose {@code turnId} is still stamped in {@code asyncTasksByTurn} —
|
||||
* not only while its question is still open. {@link #answer} deliberately leaves that stamp in
|
||||
* place ({@code clearAsyncQuestion(turnId, false)}) until the resumed turn's own future actually
|
||||
* resolves, at which point {@link #finishAsyncTask} both removes the stamp AND completes that
|
||||
* task's future in the same call. So a second task can never reach "{@code turnId} stamped, future
|
||||
* still open" — the exact pair this method matches on — while a first one already holds it: by
|
||||
* the time the stamp is gone, so is the eligibility. This is an emergent property of those two
|
||||
* facts holding together, not something this method (or its callers) enforces on its own — flip
|
||||
* {@code forgetTurn} to {@code true} in that one {@link #answer} call and it silently stops being
|
||||
* true, with nothing left to fail loudly. The callers below still handle "more than one" as
|
||||
* defence in depth against exactly that, not because they exercise it today: {@link #reply}
|
||||
* treats it as unresolvable and falls back to the inbox; {@link #abandon} would pick the oldest
|
||||
* deterministically (its own {@code matching} list has no such guarantee — see its javadoc).
|
||||
*/
|
||||
private Task askAnsweredAsyncTask(String target) {
|
||||
private List<Task> askAnsweredAsyncTasks(String target) {
|
||||
List<Task> candidates = new ArrayList<>();
|
||||
for (Task task : tasks.values()) {
|
||||
if (target.equals(task.target) && task.question == null && task.turnId != null
|
||||
&& !task.future.isDone()) {
|
||||
return task;
|
||||
candidates.add(task);
|
||||
}
|
||||
}
|
||||
return null;
|
||||
return candidates;
|
||||
}
|
||||
|
||||
/** Record a counter sample when a registry is wired; a no-op in unit tests. */
|
||||
@@ -473,16 +519,55 @@ public final class MessageService {
|
||||
* outcome is counted, so a torn-down delegation stops being invisible to {@code /metrics}.
|
||||
*
|
||||
* <p><strong>#137 defence in depth.</strong> {@link #reply} already hands a worker's real
|
||||
* {@code fleet_reply} straight to the async ticket it belongs to whenever one is still parked
|
||||
* waiting for it (see {@link #askAnsweredAsyncTask}), so by the time a session is released its
|
||||
* tasks are normally already resolved — this loop's {@code complete} calls are then harmless
|
||||
* no-ops (a {@link CompletableFuture} can only resolve once). But should some other path someday
|
||||
* strand a reply in the inbox without completing its ticket, checking
|
||||
* {@code fleet_reply} straight to the async ticket it belongs to whenever exactly one is still
|
||||
* parked waiting for it (see {@link #askAnsweredAsyncTasks}), so by the time a session is
|
||||
* released its tasks are normally already resolved — this loop's {@code complete} calls are then
|
||||
* harmless no-ops (a {@link CompletableFuture} can only resolve once). But should some other path
|
||||
* someday strand a reply in the inbox without completing its ticket, checking
|
||||
* {@link #hasStrandedReply(String)} here — before ever writing a failure — means a torn-down
|
||||
* session whose worker in fact replied is still reported {@code REPLIED} with that reply's own
|
||||
* text, never the misleading "the worker session was released before it replied" (which also
|
||||
* means the snapshot/worktree recovery hint that follows it never prints once a reply exists).
|
||||
*
|
||||
* <p><strong>At most one task gets the recovered reply — and here, unlike {@link #reply}'s
|
||||
* {@link #askAnsweredAsyncTasks}, {@code matching.size() >= 2} alone is reachable today.</strong>
|
||||
* This method's {@code matching} filter has no {@code turnId != null} requirement, so it matches
|
||||
* any plain (never-asked) open task too — and {@link #sendAsync} does not limit a target to one
|
||||
* of those: a second {@code fleet_send{wait:false}} at a target that is still busy returns its own
|
||||
* ticket immediately and simply parks its {@link #send} behind the target's session lock for up
|
||||
* to {@link #ASYNC_TIMEOUT_MS}, exactly as {@code abandonFailsEveryPendingAsyncTicketForTheReleasedTarget}
|
||||
* already proves. Before this fix, the loop below drained the strand once and then reused that
|
||||
* same {@code Reply} for <em>every</em> task it walked past — so two open tasks really did both
|
||||
* complete {@code REPLIED} with the same text (see the pre-fix loop in commit 97f6c33's parent).
|
||||
* A stranded reply is one worker answer, so it can settle at most one open task on this target —
|
||||
* never every open task, and never a guess. When more than one task is still open here, the
|
||||
* recovered reply goes to the <em>oldest</em> (lowest {@link Task#createdNanos}) — it has been
|
||||
* waiting longest, so it is the one most likely to be what the reply actually answers. Every
|
||||
* other open task keeps the ordinary {@code WORKER_FAILED} path it would take without a stranded
|
||||
* reply at all.
|
||||
*
|
||||
* <p><strong>{@code matching.size() >= 2} together with {@code hadStrandedReply} is a different
|
||||
* question, and today it is defence in depth rather than a path this codebase's public API can
|
||||
* drive.</strong> This class has exactly two sites that ever acquire a target's entry in
|
||||
* {@code sessionLocks} — {@link #send} and {@link #answer} — and both open a {@link Rendezvous}
|
||||
* waiter for that same target as the very first thing they do after acquiring the lock, then hold
|
||||
* lock and waiter together for the rest of their critical section ({@link #send} also clears
|
||||
* {@link #strandedReplies} right there, the instant it opens its waiter — before it ever enqueues
|
||||
* delivery). So "the session lock is held" and "a live waiter is open for it" are the same fact
|
||||
* throughout this class, and {@link #reply}'s fast path always resolves a currently-open waiter
|
||||
* directly rather than stranding. The two facts this method wants therefore cannot be produced
|
||||
* side by side: while the lock is held, a real reply resolves the open waiter directly and never
|
||||
* reaches {@link #strandedReplies}; the instant the lock is free, any parked matching task's own
|
||||
* {@link #send} that is scheduled next wins it and, by opening its waiter, clears the strand again
|
||||
* before this method ever runs. There is no way to hold that lock open-but-unaccepted from outside
|
||||
* {@link #send}/{@link #answer} to freeze a window in between. Constructing both facts at once
|
||||
* through {@code sendAsync}/{@code reply}/{@code ask}/{@code answer} would need a race against
|
||||
* virtual-thread scheduling, not a deterministic sequence — so the oldest-wins code below stays as
|
||||
* defence in depth against a regression to that mechanism (e.g. clearing {@link #strandedReplies}
|
||||
* on a narrower condition than "any acceptance"), not because today's test suite exercises the
|
||||
* conjunction. {@code matching.size() >= 2} alone, without a strand, is exactly what
|
||||
* {@code abandonFailsEveryPendingAsyncTicketForTheReleasedTarget} already covers.
|
||||
*
|
||||
* @return true if a live waiter or an async task was failed (never true for one recovered as a
|
||||
* reply — see the note above)
|
||||
*/
|
||||
@@ -494,21 +579,39 @@ public final class MessageService {
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
|
||||
boolean failed = waiter != null && !waiter.isDone() && rendezvous.resolveFailure(waiter, reason);
|
||||
boolean asyncFailed = false;
|
||||
Reply recovered = null; // lazily drained at most once, only if a task actually needs it
|
||||
|
||||
List<Task> matching = new ArrayList<>();
|
||||
for (Task task : tasks.values()) {
|
||||
if (!target.equals(task.target) || task.question != null || task.future.isDone()) {
|
||||
continue;
|
||||
if (target.equals(task.target) && task.question == null && !task.future.isDone()) {
|
||||
matching.add(task);
|
||||
}
|
||||
if (hadStrandedReply && recovered == null) {
|
||||
recovered = recoverStrandedReply(target);
|
||||
}
|
||||
Task recoveryTask = null;
|
||||
if (hadStrandedReply && !matching.isEmpty()) {
|
||||
recoveryTask = matching.get(0);
|
||||
for (Task candidate : matching) {
|
||||
if (candidate.createdNanos < recoveryTask.createdNanos) {
|
||||
recoveryTask = candidate;
|
||||
}
|
||||
}
|
||||
Reply outcome = recovered != null ? recovered : new Reply(Outcome.WORKER_FAILED, reason);
|
||||
}
|
||||
Reply recovered = recoveryTask != null ? recoverStrandedReply(target) : null;
|
||||
for (Task task : matching) {
|
||||
boolean isRecovery = task == recoveryTask && recovered != null;
|
||||
Reply outcome = isRecovery ? recovered : new Reply(Outcome.WORKER_FAILED, reason);
|
||||
if (task.future.complete(outcome)) {
|
||||
if (outcome.outcome() == Outcome.WORKER_FAILED) {
|
||||
asyncFailed = true;
|
||||
} else if (task.turnId != null) {
|
||||
asyncTasksByTurn.remove(task.turnId, task);
|
||||
}
|
||||
} else if (isRecovery) {
|
||||
// The recovered reply was already drained out of the inbox, but this task resolved
|
||||
// through another path (e.g. a concurrent reply() or a second abandon() racing this
|
||||
// one) between us choosing it and completing it here. Put the reply back rather than
|
||||
// lose it silently — it may still belong to some other still-open task, or the next
|
||||
// caller that drains this target's inbox.
|
||||
inbox.publish(target, UUID.randomUUID().toString(), recovered.text());
|
||||
}
|
||||
}
|
||||
if (failed) {
|
||||
@@ -523,6 +626,10 @@ public final class MessageService {
|
||||
* stranding fact raced away, e.g. a lead's own {@code fleet_poll} on the raw session already
|
||||
* drained it first). When more than one message is queued, only the newest is the worker's actual
|
||||
* final answer ({@link #drainReplies} returns them oldest-first).
|
||||
*
|
||||
* <p>This does drain (removes the messages from the inbox) before the caller knows whether the
|
||||
* task it is recovering for will actually accept them — {@link #abandon} is the one that puts a
|
||||
* reply back if its {@code complete} call turns out to lose the race.
|
||||
*/
|
||||
private Reply recoverStrandedReply(String target) {
|
||||
var messages = drainReplies(target);
|
||||
|
||||
@@ -334,8 +334,7 @@ class OpenCodeLauncherTest {
|
||||
assertNull(handle.agentSessionId(), "no record yet → null, not a spawn-time block");
|
||||
// Once the record appears (here: same cwd), lazy discovery resolves it — the handle's
|
||||
// session id matches its own worktree, not another's.
|
||||
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "p1", "ses_a.json",
|
||||
"ses_resolved", "/work/dir", 1000L);
|
||||
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_resolved", "/work/dir", 1000L);
|
||||
assertEquals("ses_resolved", handle.agentSessionId(),
|
||||
"agentSessionId() re-scans and picks up a record that has since been written");
|
||||
}
|
||||
|
||||
@@ -5,68 +5,89 @@ import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.nio.file.attribute.FileTime;
|
||||
import java.sql.Connection;
|
||||
import java.sql.DriverManager;
|
||||
import java.sql.PreparedStatement;
|
||||
import java.sql.SQLException;
|
||||
import java.sql.Statement;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
/**
|
||||
* {@link OpenCodeSessionDiscovery} matches an opencode session record by the worker's cwd (its
|
||||
* {@code directory}) against opencode's on-disk storage. These tests populate a TEMP storage root
|
||||
* themselves — never the operator's real {@code ~/.local/share/opencode}.
|
||||
* {@link OpenCodeSessionDiscovery} matches an opencode session row by the worker's cwd (its
|
||||
* {@code directory}) against opencode's {@code opencode.db} SQLite database. These tests build a
|
||||
* SYNTHETIC database themselves, in a JUnit temp directory — never the operator's real
|
||||
* {@code ~/.local/share/opencode/opencode.db}, which a live opencode process may be writing.
|
||||
*/
|
||||
class OpenCodeSessionDiscoveryTest {
|
||||
|
||||
/**
|
||||
* Write a session record {@code {"id":..., "directory":...}} under
|
||||
* {@code <root>/session/<projectID>/<fileName>} and stamp it with a known last-modified time,
|
||||
* so "most recently modified wins" is deterministic. Static so the launcher test can reuse it.
|
||||
* Create {@code <root>/opencode.db} with a minimal {@code session} table (just the columns
|
||||
* {@link OpenCodeSessionDiscovery} reads: {@code id}, {@code directory}, {@code time_updated})
|
||||
* and insert one row. Static so {@link OpenCodeLauncherTest} can reuse it.
|
||||
*/
|
||||
static void writeRecord(Path root, String projectId, String fileName, String id,
|
||||
String directory, long lastModifiedEpochMillis) throws Exception {
|
||||
Path dir = root.resolve("session").resolve(projectId);
|
||||
Files.createDirectories(dir);
|
||||
Path file = dir.resolve(fileName);
|
||||
Files.writeString(file, "{\"id\":\"" + id + "\",\"directory\":\"" + directory
|
||||
+ "\",\"projectID\":\"" + projectId + "\",\"version\":\"1.1.31\"}");
|
||||
Files.setLastModifiedTime(file, FileTime.fromMillis(lastModifiedEpochMillis));
|
||||
static void writeRecord(Path root, String id, String directory, long timeUpdated) throws Exception {
|
||||
Path db = root.resolve("opencode.db");
|
||||
try (Connection connection = DriverManager.getConnection("jdbc:sqlite:" + db)) {
|
||||
try (Statement statement = connection.createStatement()) {
|
||||
statement.execute("CREATE TABLE IF NOT EXISTS session ("
|
||||
+ "id TEXT PRIMARY KEY, directory TEXT, time_updated INTEGER)");
|
||||
}
|
||||
// Bound parameters, not string interpolation: the class under test uses a
|
||||
// PreparedStatement, and a hand-escaped INSERT here is a pattern someone copies out.
|
||||
try (PreparedStatement insert = connection.prepareStatement(
|
||||
"INSERT INTO session (id, directory, time_updated) VALUES (?, ?, ?)")) {
|
||||
insert.setString(1, id);
|
||||
insert.setString(2, directory);
|
||||
insert.setLong(3, timeUpdated);
|
||||
insert.executeUpdate();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void findsTheRecordWhoseDirectoryEqualsTheCwd(@TempDir Path root) throws Exception {
|
||||
writeRecord(root, "p1", "ses_a.json", "ses_aaa", "/w/a", 1000L);
|
||||
writeRecord(root, "p2", "ses_b.json", "ses_bbb", "/w/b", 2000L);
|
||||
void findsTheRowWhoseDirectoryEqualsTheCwd(@TempDir Path root) throws Exception {
|
||||
writeRecord(root, "ses_aaa", "/w/a", 1000L);
|
||||
writeRecord(root, "ses_bbb", "/w/b", 2000L);
|
||||
|
||||
assertEquals("ses_bbb", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/b"),
|
||||
"the record whose directory equals the cwd is the one found");
|
||||
"the row whose directory equals the cwd is the one found");
|
||||
assertEquals("ses_aaa", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aNonMatchingDirectoryYieldsNullRatherThanAMismatch(@TempDir Path root) throws Exception {
|
||||
writeRecord(root, "p1", "ses_a.json", "ses_aaa", "/w/a", 1000L);
|
||||
writeRecord(root, "ses_aaa", "/w/a", 1000L);
|
||||
|
||||
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/other"),
|
||||
"no record for this cwd yet → null, not a wrong session");
|
||||
"no row for this cwd yet → null, not a wrong session");
|
||||
}
|
||||
|
||||
@Test
|
||||
void prefersTheMostRecentlyModifiedRecordWhenSeveralMatch(@TempDir Path root) throws Exception {
|
||||
writeRecord(root, "p1", "old.json", "ses_old", "/w/a", 1000L);
|
||||
writeRecord(root, "p2", "new.json", "ses_new", "/w/a", 5000L);
|
||||
void prefersTheMostRecentlyUpdatedRowWhenSeveralMatch(@TempDir Path root) throws Exception {
|
||||
writeRecord(root, "ses_old", "/w/a", 1000L);
|
||||
writeRecord(root, "ses_new", "/w/a", 5000L);
|
||||
|
||||
assertEquals("ses_new", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"),
|
||||
"the freshest record for the cwd wins");
|
||||
"the row with the highest time_updated for the cwd wins");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aMissingOrEmptyStorageRootYieldsNullWithoutThrowing(@TempDir Path root) throws Exception {
|
||||
// Missing: no session dir at all under the root.
|
||||
void aMissingDatabaseYieldsNullWithoutThrowing(@TempDir Path root) {
|
||||
// No opencode.db at all under the root.
|
||||
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"));
|
||||
}
|
||||
|
||||
// Present but empty: a session dir with nothing in it produces no match, not a throw.
|
||||
Path emptyRoot = root.resolve("empty");
|
||||
Files.createDirectories(emptyRoot.resolve("session"));
|
||||
assertNull(new OpenCodeSessionDiscovery(emptyRoot).sessionIdForDirectory("/w/a"));
|
||||
@Test
|
||||
void anEmptyDatabaseYieldsNullWithoutThrowing(@TempDir Path root) throws Exception {
|
||||
Path db = root.resolve("opencode.db");
|
||||
try (Connection connection = DriverManager.getConnection("jdbc:sqlite:" + db);
|
||||
Statement statement = connection.createStatement()) {
|
||||
statement.execute("CREATE TABLE session (id TEXT PRIMARY KEY, directory TEXT, "
|
||||
+ "time_updated INTEGER)");
|
||||
}
|
||||
|
||||
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -76,15 +97,41 @@ class OpenCodeSessionDiscoveryTest {
|
||||
assertNull(discovery.sessionIdForDirectory(" "));
|
||||
}
|
||||
|
||||
/**
|
||||
* The one line standing between fleetd and writing the operator's live {@code opencode.db} —
|
||||
* 841MB, with a running opencode writing it — is {@code config.setReadOnly(true)} in
|
||||
* {@link OpenCodeSessionDiscovery#openReadOnly()}. Delete it and every other test in this class
|
||||
* still passes, so this is the test that guards it.
|
||||
*
|
||||
* <p>It asks the connection to write, and requires a refusal. The obvious alternative — make
|
||||
* the database file unwritable and check the read still works — proves nothing: SQLite silently
|
||||
* downgrades a read-write open of an unwritable file to read-only, so that test passes either
|
||||
* way. It was tried and watched pass with the flag removed.
|
||||
*/
|
||||
@Test
|
||||
void aMalformedRecordIsSkippedRatherThanFatal(@TempDir Path root) throws Exception {
|
||||
// A record that fails to parse must not abort the scan of its siblings.
|
||||
Path dir = root.resolve("session").resolve("p1");
|
||||
Files.createDirectories(dir);
|
||||
Files.writeString(dir.resolve("broken.json"), "{not valid json");
|
||||
writeRecord(root, "p1", "good.json", "ses_good", "/w/a", 1000L);
|
||||
void theDatabaseIsOpenedReadOnly(@TempDir Path root) throws Exception {
|
||||
writeRecord(root, "ses_aaa", "/w/a", 1000L);
|
||||
|
||||
assertEquals("ses_good", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"),
|
||||
"an unreadable record is skipped; a later valid one still matches");
|
||||
try (Connection connection = new OpenCodeSessionDiscovery(root).openReadOnly();
|
||||
Statement statement = connection.createStatement()) {
|
||||
SQLException refused = assertThrows(SQLException.class,
|
||||
() -> statement.executeUpdate("INSERT INTO session (id, directory, time_updated) "
|
||||
+ "VALUES ('ses_zzz', '/w/z', 1)"),
|
||||
"the connection must REFUSE a write — opencode is writing this database live");
|
||||
assertTrue(refused.getMessage().toLowerCase().contains("readonly")
|
||||
|| refused.getMessage().toLowerCase().contains("read-only"),
|
||||
"the refusal must be about read-only, not some other error: " + refused.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void aCorruptDatabaseFileYieldsNullWithoutThrowing(@TempDir Path root) throws Exception {
|
||||
// A file at opencode.db that is not a SQLite database at all — the open/query must fail
|
||||
// safe, never fatal to a spawn.
|
||||
Path db = root.resolve("opencode.db");
|
||||
Files.writeString(db, "this is not a sqlite database");
|
||||
|
||||
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"),
|
||||
"an unreadable database resolves to null, not an exception");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -688,6 +688,32 @@ class MessageServiceTest {
|
||||
assertFailedTicket(third, "agent target term_a not found");
|
||||
}
|
||||
|
||||
// --- #137 follow-up: abandon() must not guess when more than one task is open ---------------
|
||||
//
|
||||
// A test combining a genuine stranded reply (hasStrandedReply(T)==true) with two simultaneously
|
||||
// open matching tasks was attempted here and removed after investigation showed the combination
|
||||
// is not reachable through the public API today, not merely hard to time right:
|
||||
//
|
||||
// This class has exactly two call sites that ever hold a target's entry in the session-lock map
|
||||
// (send() and answer()), and both open a Rendezvous waiter for that same target as the first thing
|
||||
// they do after acquiring the lock, holding lock and waiter together for their whole critical
|
||||
// section. So "the lock is held" and "a live waiter is open" are the same fact throughout this
|
||||
// class. reply()'s fast path always resolves a currently-open waiter directly instead of
|
||||
// stranding — so a strand can only be created while NO task is accepted (lock free), and the
|
||||
// instant the lock is next taken (by any parked matching task's own send(), the moment it is
|
||||
// scheduled), that acceptance clears strandedReplies again (see send()'s CB-640 comment) before
|
||||
// abandon() can ever observe both facts together. Confirmed empirically too: an earlier version of
|
||||
// this test stranded a reply, then created an "accepted" task (awaitWaiting()) followed by a
|
||||
// "parked" one — and the accepted task's own acceptance silently cleared the strand it was
|
||||
// supposed to be racing against, so the parked task came back WORKER_FAILED instead of DONE, not
|
||||
// because the fix was missing but because the test's premise could not be constructed.
|
||||
//
|
||||
// The reachable half — matching.size() >= 2 alone, no strand — is exactly what
|
||||
// abandonFailsEveryPendingAsyncTicketForTheReleasedTarget already covers (all fail, none guess).
|
||||
// The oldest-wins code in abandon() stays as defence in depth (see its own javadoc) against a
|
||||
// regression that would make the conjunction reachable, e.g. clearing strandedReplies on a
|
||||
// narrower condition than "any acceptance" — not because this suite exercises it today.
|
||||
|
||||
@Test
|
||||
void abandonDoesNotFailAnAsyncTicketWaitingForAnAnswer() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
|
||||
Reference in New Issue
Block a user