commit - 31d2568cf13643602ba7bc0dd80bfb086b763bdc
commit + bc50b132dba26f836a2a47e083934d55aa50d9f9
blob - a5983bde67263c1d8fe2f3a5033646e035903189
blob + 98e6139ae3aaffc0f149947dcdf6ed08e1dc99ab
--- regress/turn/run.sh
+++ regress/turn/run.sh
fi
stop_stub
+# A ROTATE frame is ownership state in its own right. Publish it with a
+# blocking selftest send, then cancel before the ordinary Tool result. The
+# background job is still pending in this Turn, so rollback must consume the
+# frame without trying to release committed job ownership.
+rm -rf "$HOME/.fugu"
+rm -f "$dir/out" "$dir/errout" "$dir/reqout"
+start_stub "BGTOOL JOBOUTPUT"
+write_conf "$stubport"
+FUGU_SELFTEST=1 FUGU_ARTIFACT_FRAME0_DELAY_MS=3000 \
+ FUGU_CONF=$dir/fugu.conf FUGU_LIBEXEC=$dir/libexec/fugu \
+ FUGU_CA_FILE=$dir/cert.pem FUGU_RETRY_BASE_MS=50 \
+ "$fugu" -d -p "cancel pending Artifact rotation" </dev/null \
+ >"$dir/out" 2>"$dir/errout" &
+fugupid=$!
+i=0
+while ! grep -Fq 'selftest frame0 Artifact rotation published' \
+ "$dir/errout" 2>/dev/null && [ "$i" -lt 200 ] && \
+ kill -0 "$fugupid" 2>/dev/null; do
+ sleep 0.05
+ i=$((i + 1))
+done
+grep -Fq 'selftest frame0 Artifact rotation published' \
+ "$dir/errout" 2>/dev/null
+reached=$?
+if [ "$reached" -eq 0 ]; then
+ /bin/kill -INT "$fugupid" 2>/dev/null
+else
+ /bin/kill -TERM "$fugupid" 2>/dev/null
+fi
+wait "$fugupid"
+rc=$?
+fugupid=''
+ok "background rotation cancel: pending frame publishes before Ctrl-C" \
+ "$reached"
+ok "background rotation cancel: pending Turn exits 3" \
+ $([ "$rc" -eq 3 ]; echo $?)
+n=$(grep -h -c '"outcome":"abandoned"' \
+ "$sessions"/*.ndjson 2>/dev/null | awk '{ n += $1 } END { print n + 0 }')
+ok "background rotation cancel: pending Turn is abandoned once" \
+ $([ "$n" -eq 1 ]; echo $?)
+n=$(grep -h -c '"outcome":"ok"' "$sessions"/*.ndjson 2>/dev/null | \
+ awk '{ n += $1 } END { print n + 0 }')
+ok "background rotation cancel: pending job never commits" \
+ $([ "$n" -eq 0 ]; echo $?)
+n=$(grep -Fc 'selftest frame0 Artifact rotation published' \
+ "$dir/errout" 2>/dev/null)
+ok "background rotation cancel: pending frame publishes exactly once" \
+ $([ "$n" -eq 1 ]; echo $?)
+if grep -Eq 'selftest Artifact RELEASE|bad Artifact|mismatched Artifact|stale Artifact|missing cancelled Artifact' \
+ "$dir/errout" 2>/dev/null; then
+ fail "background rotation cancel: pending ownership rolls back cleanly"
+else
+ pass "background rotation cancel: pending ownership rolls back cleanly"
+fi
+stop_stub
+
+# The same cancellation after the launch Turn commits must preserve the
+# background job but release the frozen side of the published rotation.
+# Curses routes Ctrl-C through the normal UI cancellation handshake and exits
+# only after Coordinator has settled that ownership frame.
+if [ -x "$obj/ptydrive" ]; then
+ rm -rf "$HOME/.fugu"
+ rm -f "$dir/reqout" "$dir/reqlog" "$dir/rotation-cursesout"
+ start_stub "BGTOOL OK JOBOUTPUT"
+ write_conf "$stubport"
+ printf 'launch committed background\r\000\000\000poll committed background\r\000\003' | \
+ TERM=vt100 HOME=$HOME FUGU_SELFTEST=1 \
+ FUGU_ARTIFACT_FRAME0_DELAY_MS=3000 \
+ FUGU_CONF=$dir/fugu.conf FUGU_LIBEXEC=$dir/libexec/fugu \
+ FUGU_CA_FILE=$dir/cert.pem FUGU_RETRY_BASE_MS=50 \
+ "$obj/ptydrive" "$fugu" -d \
+ >"$dir/rotation-cursesout" 2>&1
+ rc=$?
+ ok "background rotation cancel: curses exits after Ctrl-C" \
+ $([ "$rc" -eq 0 ]; echo $?)
+ n=$(grep -h -c '"outcome":"ok"' "$sessions"/*.ndjson 2>/dev/null | \
+ awk '{ n += $1 } END { print n + 0 }')
+ ok "background rotation cancel: launch Turn remains committed" \
+ $([ "$n" -eq 1 ]; echo $?)
+ n=$(grep -h -c '"outcome":"abandoned"' \
+ "$sessions"/*.ndjson 2>/dev/null | \
+ awk '{ n += $1 } END { print n + 0 }')
+ ok "background rotation cancel: poll Turn is abandoned once" \
+ $([ "$n" -eq 1 ]; echo $?)
+ n=$(grep -Fc 'selftest frame0 Artifact rotation published' \
+ "$dir/rotation-cursesout" 2>/dev/null)
+ ok "background rotation cancel: committed frame publishes once" \
+ $([ "$n" -eq 1 ]; echo $?)
+ n=$(grep -Fc 'selftest Artifact RELEASE job=1' \
+ "$dir/rotation-cursesout" 2>/dev/null)
+ ok "background rotation cancel: committed frozen spool releases once" \
+ $([ "$n" -eq 1 ]; echo $?)
+ if grep -Eq 'selftest Artifact CANCEL_ACK|bad Artifact|mismatched Artifact|stale Artifact|missing cancelled Artifact' \
+ "$dir/rotation-cursesout" 2>/dev/null; then
+ fail "background rotation cancel: committed ownership settles cleanly"
+ else
+ pass "background rotation cancel: committed ownership settles cleanly"
+ fi
+ stop_stub
+else
+ fail "background rotation cancel: ptydrive helper not built"
+fi
+
+# Force a completed committed job into normal archival while its Artifact
+# terminal is delayed and converted to STORE_ERROR. DROP is already queued
+# on the Tool socket when Coordinator sends CANCEL; the ordered Tool output
+# must therefore be DROP followed by the exact job-wide CANCEL_ACK.
+rm -rf "$HOME/.fugu"
+rm -f "$dir/out" "$dir/errout" "$dir/reqout"
+start_stub "BGFAST OK RETRY429 OK JOBOUTPUT OK"
+write_conf "$stubport"
+printf 'launch short background\nwait for store failure\nprobe cancelled background\n/quit\n' \
+ > "$dir/lin"
+FUGU_SELFTEST=1 FUGU_ARTIFACT_ARCHIVE_COMPLETED=1 \
+ FUGU_ARTIFACT_DROP_DELAY_MS=1500 \
+ FUGU_ARTIFACT_STORE_ERROR_DELAY_MS=500 \
+ FUGU_CONF=$dir/fugu.conf FUGU_LIBEXEC=$dir/libexec/fugu \
+ FUGU_CA_FILE=$dir/cert.pem FUGU_RETRY_BASE_MS=50 \
+ "$fugu" -d <"$dir/lin" >"$dir/out" 2>"$dir/errout"
+rc=$?
+ok "background retirement race: session completes" \
+ $([ "$rc" -eq 0 ]; echo $?)
+n=$(grep -h -c '"outcome":"ok"' "$sessions"/*.ndjson 2>/dev/null | \
+ awk '{ n += $1 } END { print n + 0 }')
+ok "background retirement race: all three Turns commit" \
+ $([ "$n" -eq 3 ]; echo $?)
+grep -q 'shell_output: no such job 1' "$dir/reqout" 2>/dev/null
+ok "background retirement race: failed job ownership is gone" $?
+drop_line=$(grep -n -m 1 'selftest Artifact DROP queued' \
+ "$dir/errout" 2>/dev/null | cut -d: -f1)
+error_line=$(grep -n -m 1 'selftest STORE_ERROR injected' \
+ "$dir/errout" 2>/dev/null | cut -d: -f1)
+ack_line=$(grep -n -m 1 'selftest Artifact CANCEL_ACK job=1' \
+ "$dir/errout" 2>/dev/null | cut -d: -f1)
+[ -n "$drop_line" ] && [ -n "$error_line" ] && [ -n "$ack_line" ] && \
+ [ "$drop_line" -lt "$error_line" ] && [ "$error_line" -lt "$ack_line" ]
+ok "background retirement race: DROP precedes STORE_ERROR and CANCEL_ACK" $?
+n=$(grep -Ec 'selftest Artifact DROP queued|selftest STORE_ERROR injected|selftest Artifact CANCEL_ACK job=1' \
+ "$dir/errout" 2>/dev/null)
+ok "background retirement race: each protocol trace occurs exactly once" \
+ $([ "$n" -eq 3 ]; echo $?)
+if grep -Eq 'bad Artifact|mismatched Artifact|stale Artifact|unexpected Artifact|worker .* exited' \
+ "$dir/errout" 2>/dev/null; then
+ fail "background retirement race: exact-key composition stays valid"
+else
+ pass "background retirement race: exact-key composition stays valid"
+fi
+stop_stub
+
echo "$checks checks, $fails failures"
[ "$fails" -eq 0 ]
blob - fbccbbe31d6cad6698801d3e6f1ba7f266ad0810
blob + 98050bde55120ecfd0bf3740de52ea3a9dcadcb2
--- regress/turn/stub.c
+++ regress/turn/stub.c
* WEBFETCH (a brokered web_fetch), OUTPUTREAD (an explicit brokered Artifact
* page), OUTPUTREADDEFAULT (a defaults-only Artifact page),
* BGTOOL (shell_bg sleep 300), BGSHORT (shell_bg sleep 1),
+ * BGFAST (shell_bg immediate output),
* ERROR_BG (the provider failure after
* that launch), JOBOUTPUT (shell_output job 1), DOUBLE (a hostile error event
* appended after a complete OK stream), AGENTCALL (lead delegates to a
"event: message_stop\n"
"data: {\"type\":\"message_stop\"}\n\n";
+/* Complete comfortably before the retrying next Turn so its current spool
+ * can be archived without a model-requested shell_output poll. */
+static const char sse_bgfast[] =
+"event: message_start\n"
+"data: {\"type\":\"message_start\",\"message\":{\"usage\":"
+ "{\"input_tokens\":9,\"output_tokens\":1}}}\n\n"
+"event: content_block_start\n"
+"data: {\"type\":\"content_block_start\",\"index\":0,"
+ "\"content_block\":{\"type\":\"tool_use\",\"id\":\"toolu_fast\","
+ "\"name\":\"shell_bg\"}}\n\n"
+"event: content_block_delta\n"
+"data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":"
+ "{\"type\":\"input_json_delta\",\"partial_json\":"
+ "\"{\\\"command\\\":\\\"echo fast-done\\\"}\"}}\n\n"
+"event: content_block_stop\n"
+"data: {\"type\":\"content_block_stop\",\"index\":0}\n\n"
+"event: message_delta\n"
+"data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":"
+ "\"tool_use\"},\"usage\":{\"output_tokens\":12}}\n\n"
+"event: message_stop\n"
+"data: {\"type\":\"message_stop\"}\n\n";
+
/* Probe the exact id allocated by the preceding BGTOOL response. If the
* failed turn was rolled back, shell_output returns "no such job 1" and
* that tool error appears in the following request. */
write_chunk(conn, sse_bgshort,
strlen(sse_bgshort));
(void)tls_write_all(conn, "0\r\n\r\n", 5);
+ } else if (strcmp(argv[i], "BGFAST") == 0) {
+ (void)tls_write_all(conn, body_ok,
+ strlen(body_ok));
+ write_chunk(conn, sse_bgfast,
+ strlen(sse_bgfast));
+ (void)tls_write_all(conn, "0\r\n\r\n", 5);
} else if (strcmp(argv[i], "ERROR_BG") == 0) {
(void)tls_write_all(conn, body_ok,
strlen(body_ok));
blob - 31bf64155a0b478d0b4f6aaf0b53133f2934f272
blob + 99c942ff5464de8407bbe7eb2c5268172ce75424
--- src/fugu-artifact/main.c
+++ src/fugu-artifact/main.c
static struct artifact_state *cleanup_state;
static void
+artifact_store_error_test(struct artifact_terminal *terminal)
+{
+ static int injected;
+ const char *value, *errstr;
+ struct timespec delay;
+ long long ms;
+
+ if (injected || issetugid() || getenv("FUGU_SELFTEST") == NULL ||
+ (value = getenv("FUGU_ARTIFACT_STORE_ERROR_DELAY_MS")) == NULL)
+ return;
+ ms = strtonum(value, 1, 5000, &errstr);
+ if (errstr != NULL)
+ return;
+ injected = 1;
+ delay.tv_sec = (time_t)(ms / 1000);
+ delay.tv_nsec = (long)(ms % 1000) * 1000000L;
+ (void)nanosleep(&delay, NULL);
+ log_warnx("fugu-artifact: selftest STORE_ERROR injected");
+ terminal->status = ARTIFACT_STORE_ERROR;
+}
+
+static void
artifact_read_test_delay(void)
{
const char *value, *errstr;
rc = artifact_pending_drain(de->pending, &terminal);
if (rc == 0)
return;
+ artifact_store_error_test(&terminal);
event_del(&de->ev);
drain_unlink(de);
imsgev_send(&de->state->worker->iev, FUGU_IMSG_ARTIFACT_TERMINAL,
blob - 690148883bf6750cb315d158868210944f9fb89a
blob + ba537a15eb22ee38f21810488368ceb153bbbf93
--- src/fugu-tool/job.c
+++ src/fugu-tool/job.c
jl->record_limit = (size_t)cap + JOB_COMPLETED_RECORD_MAX;
}
+void
+job_selftest_set_ownership_limit(struct joblist *jl, size_t limit)
+{
+ if (jl == NULL || limit == 0 || limit > JOB_OWNERSHIP_MAX)
+ abort();
+ jl->ownership_limit = limit;
+}
+
static struct job *
job_find(struct joblist *jl, int id)
{
(*oldest_link)->artifact_managed));
}
+int
+job_selftest_archive_completed(struct joblist *jl)
+{
+ struct job *j;
+
+ if (jl == NULL)
+ abort();
+ for (j = jl->head; j != NULL; j = j->next) {
+ if (j->pending || j->artifact_settling || j->fd != -1 ||
+ j->artifact_fd != -1 || j->pending_write.len != 0)
+ continue;
+ job_stop(j);
+ break;
+ }
+ return (job_prune_completed_one(jl));
+}
+
static int
job_make_record_room(struct joblist *jl)
{
blob - 796e2ecc1bb8cb70f469ead0ccee52e97306427a
blob + 8e3b9e77a89f6c8b6d7171161b8dbd1721937fad
--- src/fugu-tool/job.h
+++ src/fugu-tool/job.h
/* the max_bg_jobs gate; 0 (the default here) means unconfigured */
void job_set_cap(struct joblist *, int cap);
+/* Regression-only bound selected by fugu-tool under guarded selftest mode. */
+void job_selftest_set_ownership_limit(struct joblist *, size_t);
+int job_selftest_archive_completed(struct joblist *);
void job_start(struct joblist *, struct tool_ctx *, const struct json *,
int obj, struct buf *out, int *is_error);
blob - 89929b783364ee5513ff07fe97e5a87e7337572c
blob + 20e8d5ffe78a9765287b58c24e371c50dba899ef
--- src/fugu-tool/main.c
+++ src/fugu-tool/main.c
#include <stdlib.h>
#include <string.h>
#include <syslog.h>
+#include <time.h>
#include <unistd.h>
#include "log.h"
#include "artifact_scope.h"
static int readonly; /* -r: the subagent role, always read-only */
+static long long selftest_rotation_delay_ms;
+static long long selftest_drop_delay_ms;
+static int selftest_archive_pending;
+static int selftest_rotation_trace;
+static int selftest_drop_trace;
struct tool_state {
struct worker *worker;
* a clean shutdown that would otherwise orphan them into init. */
static struct joblist *jobs_atexit;
+static long long
+selftest_delay(const char *name)
+{
+ const char *value, *errstr;
+ long long ms;
+
+ if (issetugid() || getenv("FUGU_SELFTEST") == NULL ||
+ (value = getenv(name)) == NULL)
+ return (0);
+ ms = strtonum(value, 1, 5000, &errstr);
+ if (errstr != NULL)
+ fatalx("fugu-tool: invalid %s", name);
+ return (ms);
+}
+
static int
+selftest_flag(const char *name)
+{
+ const char *value;
+
+ if (issetugid() || getenv("FUGU_SELFTEST") == NULL ||
+ (value = getenv(name)) == NULL)
+ return (0);
+ if (strcmp(value, "1") != 0)
+ fatalx("fugu-tool: invalid %s", name);
+ return (1);
+}
+
+static void
+selftest_pause(long long *milliseconds, const char *marker)
+{
+ struct timespec delay;
+ long long ms;
+
+ if ((ms = *milliseconds) == 0)
+ return;
+ *milliseconds = 0;
+ log_warnx("%s", marker);
+ delay.tv_sec = (time_t)(ms / 1000);
+ delay.tv_nsec = (long)(ms % 1000) * 1000000L;
+ (void)nanosleep(&delay, NULL);
+}
+
+static int
spool_drop_id_valid(const char id[FUGU_ARTIFACT_ID_BUFSZ])
{
static const char zero[FUGU_ARTIFACT_ID_BUFSZ];
fatalx("fugu-tool: invalid Artifact spool drop");
imsgev_send(&w->iev, FUGU_IMSG_TOOL_ARTIFACT_DROP, correlation,
&drop, sizeof(drop));
+ selftest_pause(&selftest_drop_delay_ms,
+ "fugu-tool: selftest Artifact DROP queued");
}
}
(void)ev;
if (ts->ctx != NULL) {
job_pump(ts->ctx->jobs);
+ if (selftest_archive_pending &&
+ job_selftest_archive_completed(ts->ctx->jobs) == 1)
+ selftest_archive_pending = 0;
send_spool_drops(ts->worker, ts);
}
}
ts->ctx = tool_ctx_new(write_ok, net_ok);
ts->worker = w;
job_set_cap(ts->ctx->jobs, (int)ts->cap);
+ if (selftest_rotation_delay_ms > 0)
+ job_selftest_set_ownership_limit(ts->ctx->jobs, 1);
/* kill our background jobs when we exit, and keep their output
* draining between polls (event_init ran before this dispatch) */
if (ts->ctx->artifact_spool_adopted && ts->ctx->artifact_available)
strlcpy(rotate.replacement_id, ts->artifact.id,
sizeof(rotate.replacement_id));
+ if (selftest_rotation_delay_ms > 0) {
+ imsgev_send_blocking(&w->iev, FUGU_IMSG_TOOL_ARTIFACT_ROTATE, id,
+ &rotate, sizeof(rotate));
+ selftest_pause(&selftest_rotation_delay_ms,
+ "fugu-tool: selftest frame0 Artifact rotation published");
+ return;
+ }
imsgev_send(&w->iev, FUGU_IMSG_TOOL_ARTIFACT_ROTATE, id, &rotate,
sizeof(rotate));
}
sizeof(ack.artifact_id));
imsgev_send(&w->iev, FUGU_IMSG_TOOL_ARTIFACT_CANCEL_ACK,
job_id, &ack, sizeof(ack));
+ if (selftest_drop_trace)
+ log_warnx("fugu-tool: selftest Artifact CANCEL_ACK "
+ "job=%u", job_id);
+ } else if (selftest_rotation_trace) {
+ log_warnx("fugu-tool: selftest Artifact RELEASE job=%u",
+ job_id);
}
return (0);
}
buf_init(&ts.arg);
ts.artifact_fd = -1;
agent_boot_init(&ts.boot);
+ selftest_rotation_delay_ms = selftest_delay(
+ "FUGU_ARTIFACT_FRAME0_DELAY_MS");
+ selftest_drop_delay_ms = selftest_delay("FUGU_ARTIFACT_DROP_DELAY_MS");
+ selftest_archive_pending = selftest_flag(
+ "FUGU_ARTIFACT_ARCHIVE_COMPLETED");
+ selftest_rotation_trace = selftest_rotation_delay_ms > 0;
+ selftest_drop_trace = selftest_drop_delay_ms > 0;
/*
* Transitional pledge. It is a superset of the steady promises;