Inbound Message Reliability
A sender can receive no reply even though the agent started the work. A restart can also cause work or an external action to happen again. Before resubmitting a silent request, the operator must inspect the session, outbox, bridge outcome and any action the request could have caused.
Redis recovery retries delivery to dispatch, not to successful completion. The agent queues acknowledgment when it hands a message to the agent loop, without waiting for an answer. Once that acknowledgment lands, Redis will not recover an interrupted turn. If it does not land, reclaim can dispatch the message again. There is no exactly-once guarantee for completion or effects.
Session recovery is separate. On restart, the agent resumes its saved session, and an idle reprompt can act on the unanswered request again. The active message and reply-source binding lived only in process memory. The work can recur while the reply loses its source binding: a resumed answer may reach the session but no source outbox, leaving the original sender with silence.
| Observation | What it establishes |
|---|---|
| Service is running | The process is up; useful progress still needs inspection |
| No pending Redis entries | No entries await acknowledgment in that group |
| Acknowledgment landed | Dispatch was acknowledged, not completion |
| Reply in the agent outbox | A reply was queued for the bridge |
| Session resumed | Conversation state returned; source routing may not have |
There is no durable deduplication across agent restarts and no durable binding that reconnects a resumed turn to its source. Within a process, a bounded seen window suppresses repeat dispatch. Failed acknowledgments retry the acknowledgment first; exhausting that budget allows reclaim to offer the request again. External effects are not transactional with acknowledgment.
The operator procedure
Use the host-side operator identity to inspect missing or duplicate replies.
Check the agent with cowboy status and cowboy doctor first. A running
service does not prove that a request completed or reached its sender.
The identity. Not the agent’s. Every agent credential is scoped to
~<agent>:* and holds only the verbs the poll issues, which does not include
XPENDING — the agent has no use for it, and an agent that could enumerate its
own pending list still could not see another’s. modules/pubsub.nix generates
one identity for a person instead: operator, readable across the whole
keyspace, with no @admin (it cannot rewrite the ACL) and no @dangerous (no
CONFIG, FLUSHALL, KEYS, SHUTDOWN). Its URL is root-only, and Redis
binds inside the namespace, so every command below is ip netns exec:
sudo -i
. /run/cowboy/operator-redis.env # sets REDIS_URL
r() { ip netns exec cowboy-ns redis-cli -u "$REDIS_URL" "$@"; }
Three names come from the generated configuration rather than from memory. For
agent <agent> and source <source>, /etc/cowboy/<agent>/sources.json names
the stream; the group is cowboy; the consumer is harness-<source> — the
SOURCE, not the agent (RedisStreamsSource::new names it after the source it
polls; two agents on the same source carry the same consumer name on their
own separate streams):
stream=$(jq -r '.sources["<source>"].stream' /etc/cowboy/<agent>/sources.json)
group=cowboy
consumer=harness-<source>
1. Is anything stuck? The pending list is the only record of an entry the agent read and did not acknowledge:
r XPENDING "$stream" "$group"
A summary of 0 means the group has no pending entries. It does not prove
that every stream entry was read or that dispatched requests were answered.
2. Which entry, and how long has it been stuck? The extended form lists one line per entry — id, consumer, idle milliseconds, delivery count:
r XPENDING "$stream" "$group" - + 10
An entry whose idle time is under ten minutes is not yet eligible for reclaim.
It may be queued, or its reader may have stopped; age alone does not prove
progress. Reclaim waits until RECLAIM_IDLE_MS has passed. An entry idle for longer
than that whose delivery count is not rising is one no poll is reaching — check
that the agent is running at all before doing anything else.
3. What was in it?
r XRANGE "$stream" <id> <id>
4. Did the turn already have an effect? Ask this before step 5, because step 5 is the one action in this procedure that can do harm. The agent’s reply goes to its own outbox, and that stream is the record of it:
r XRANGE "<agent>:<source>:outbox" - + COUNT 20
An outbox entry shows a reply was queued, not that the external platform delivered it. Inspect bridge delivery state and the destination before resubmitting. Also check files, rebuild results and other requested actions.
Do not assume reclaim will discard a request that already had an effect. Deduplication is bounded and lives in process memory; after restart the same entry can dispatch again. Neither reclaim nor manual resubmission is an exactly-once operation. If work already happened, reconcile that outcome rather than blindly replaying the original request.
5. Resubmit only after reconciliation. If retrying is appropriate, submit a new entry. Keep in mind that the old pending entry can still be reclaimed. This minimal command submits content only; it does not restore the original sender or channel metadata and therefore does not promise a routed reply:
r XADD "$stream" '*' content '<the content from step 3>'
Implementation
What one poll does
poll_command emits three redis-cli calls in one sh -c:
XGROUP CREATE … 0-0 MKSTREAM, idempotent, so a stream that does not exist yet is not an error.XAUTOCLAIM <stream> <group> <consumer> <idle> <cursor> COUNT 10— the recovery step.XREADGROUP >moves an entry into the group’s Pending Entries List under the consumer that read it, and never offers a pending entry again —>means “entries no consumer has read”. So an entry a crashed run had taken is invisible to every subsequent>, and the reclaim is the only thing that hands it back. The consumer name is deliberately stable across restarts (RedisStreamsSourcedefaults it toharness-<source>, and the comment there says why): a fresh name per process would leave the whole previous pending list keyed to a consumer that will never ask for anything again. Stability keeps the list reachable; the reclaim is what reaches it.XREADGROUP GROUP … COUNT 10 STREAMS <stream> '>'— new entries.
Reclaim before read, so a message stranded by a restart is not queued behind
messages that just arrived. A marker line separates the two replies; both are
parsed by the same hand-rolled reader over redis-cli --no-raw output.
The scan carries its cursor. XAUTOCLAIM scans the pending list from a
cursor and answers with the one to resume from. A scan that spends its COUNT
on entries too young to claim returns an empty batch and a non-zero cursor;
restarting from 0-0 every poll re-walks that same head, and an orphan behind
a full scan’s worth of young entries is only reached once the head drains on
its own. RedisStreamsSource::reclaim_cursor keeps the returned cursor, which
is why poll_command and parse_poll_result take &mut self. Redis answers
0-0 when the scan wraps, so the walk returns to the head on its own.
An entry with no usable content is still emitted, with an empty one. It
has to be: an entry nothing returns is an entry nothing acknowledges, and since
the poll reclaims the pending list it would come back on every poll forever,
masking the ten entries behind it. The manager acknowledges and drops it.
Acknowledgment semantics
XACK answers with the number of entries it removed from the pending list.
The reply is examined, not the exit code alone: redis-cli prints a refused
command’s error text on stdout and exits zero, so an ACL that does not grant
XACK used to read as an acknowledgment that landed. The generated command
requires the reply to be a run of digits — zero included, which is the
idempotent case — and exits ACK_REPLY_NOT_A_COUNT otherwise.
The acknowledgment is queued at dispatch, not at completion.
process_pubsub_messages calls SourceManager::mark_processed on each message
as soon as the prompt has been handed to process_user_input, and the same
poll tick then calls send_pubsub_acks. Nothing waits for the turn to produce
an answer. Redis’s redelivery therefore covers exactly one window — from the
XREADGROUP that read the entry to the XACK that follows its dispatch — and
nothing after it. A process that dies mid-turn has already acknowledged the
entry it was working on; the pending list holds no record of it, and no reclaim
will bring it back. That is the deliberate trade: the alternative, acking at
completion, makes every crashed turn a guaranteed re-answer, including the ones
that had already sent their reply.
A failed acknowledgment retries the acknowledgment, not the request.
SourceManager::retry_ack queues the same XACK again, up to
MAX_ACK_ATTEMPTS. Only once that budget is spent does the manager forget the
entry, which hands it back to the reclaim. Forgetting on the first failure —
before the retry budget is spent — gives the reclaim a message the agent
has already answered, and the next poll starts a second turn on it.