Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,8 @@

## Unreleased

- Sender: retain native `reply_contexts` in delivery feedback after inbound-cache
eviction and held-reply batching, so hosts can acknowledge only completed turns.
- Sender: add the host-only `handle_agent_reply/4` callback for immutable turn
context, keeping delayed replies in their original conversation and parent
after slot reuse without granting the rebound slot edit authority. The reply
Expand Down
7 changes: 7 additions & 0 deletions lib/genswarms/telegram/delivery_effects.ex
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,13 @@ defmodule Genswarms.Telegram.DeliveryEffects do
Adapters may be configured either as `Module` or `{Module, opts}`. Tuple
adapters can implement the callback with one extra final `opts` argument; the
package will prefer that arity when it exists.

For text replies, `after_delivery` metadata includes `:reply_contexts`, a list
of original native `handle_agent_reply/4` contexts (empty for ordinary messages).
Held replies retain all distinct contexts of the texts combined in that send.
Hosts may acknowledge those turns only when the outcome is successful. These
contexts survive cache eviction; `:reply_to_message_id` remains the independently
validated Telegram tag and can be nil. Model message fields cannot set contexts.
"""

@callback before_send(map()) :: :ok | {:error, term()}
Expand Down
65 changes: 46 additions & 19 deletions lib/genswarms/telegram/objects/sender.ex
Original file line number Diff line number Diff line change
Expand Up @@ -238,7 +238,7 @@ defmodule Genswarms.Telegram.Objects.Sender do
end

if valid_cid?(cid) and authorized? do
{:ok, state} = send_text(from, cid, msg, state, :reply)
{:ok, state} = send_text(from, cid, msg, state, :reply, [context])
{:noreply, state}
else
{:noreply, state}
Expand Down Expand Up @@ -1339,7 +1339,16 @@ defmodule Genswarms.Telegram.Objects.Sender do
end
end

defp send_text(from, cid, msg, state, origin) do
defp send_text(from, cid, msg, state, origin, reply_contexts \\ []) do
# Only the native callback supplies these contexts. Model message fields
# cannot create them, and the bounded inbound cache cannot invalidate them.
meta = %{
origin: origin,
from: from,
mark: Map.get(msg, "mark"),
reply_contexts: reply_contexts
}

with {:cont, state} <- prepare_delivery(from, cid, origin, state) do
text =
Adapter.call(state.delivery_effects, :redact_outbound, [
Expand All @@ -1354,29 +1363,20 @@ defmodule Genswarms.Telegram.Objects.Sender do

state =
state
|> record_logical_delivery(cid, %{text: text}, result, %{
origin: origin,
from: from,
text: text,
mark: Map.get(msg, "mark")
})
|> record_logical_delivery(cid, %{text: text}, result, Map.put(meta, :text, text))
|> stamp_reply(cid, origin)

{:ok, state}
else
case do_send_text(cid, text, msg, state, %{
origin: origin,
from: from,
mark: Map.get(msg, "mark")
}) do
case do_send_text(cid, text, msg, state, meta) do
{:ok, state} -> {:ok, state |> stamp_reply(cid, origin) |> stamp_sig(cid, text, origin)}
other -> other
end
end
else
{:suppress, cid, state} ->
parent = validate_reply_tag(cid, Map.get(msg, "reply_to_message_id"), state)
{:ok, hold_reply(from, cid, Map.get(msg, "text", ""), state, parent)}
{:ok, hold_reply(from, cid, Map.get(msg, "text", ""), state, parent, reply_contexts)}
end
end

Expand Down Expand Up @@ -3707,7 +3707,7 @@ defmodule Genswarms.Telegram.Objects.Sender do
# message when the window expires (edits don't notify on Telegram, so
# append-by-edit would deliver the answer silently). Exact replays of the
# just-delivered text — the original spam case — still die.
defp hold_reply(from, cid, text, state, parent \\ nil) do
defp hold_reply(from, cid, text, state, parent \\ nil, reply_contexts \\ []) do
text = String.trim(to_string(text))
cur = Map.get(state.held, cid)
held_len = if cur, do: cur.texts |> Enum.map(&String.length/1) |> Enum.sum(), else: 0
Expand All @@ -3720,7 +3720,14 @@ defmodule Genswarms.Telegram.Objects.Sender do
state

cur != nil and text in cur.texts ->
state
update_in(state.held[cid], fn h ->
parent = if h.from == from and Map.get(h, :reply_to) == parent, do: parent

Map.merge(h, %{
reply_to: parent,
reply_contexts: Enum.uniq(Map.get(h, :reply_contexts, []) ++ reply_contexts)
})
end)

cur == nil and map_size(state.held) >= @held_cids_max ->
state
Expand All @@ -3734,11 +3741,24 @@ defmodule Genswarms.Telegram.Objects.Sender do
true ->
state = if cur == nil, do: schedule_held_flush(cid, state), else: state

entry = %{
texts: [text],
from: from,
reply_to: parent,
reply_contexts: reply_contexts
}

held =
Map.update(state.held, cid, %{texts: [text], from: from, reply_to: parent}, fn h ->
Map.update(state.held, cid, entry, fn h ->
# A combined tail may only identify a parent shared by every text.
parent = if h.from == from and Map.get(h, :reply_to) == parent, do: parent
Map.merge(h, %{texts: h.texts ++ [text], from: from, reply_to: parent})

Map.merge(h, %{
texts: h.texts ++ [text],
from: from,
reply_to: parent,
reply_contexts: Enum.uniq(Map.get(h, :reply_contexts, []) ++ reply_contexts)
})
end)

%{state | held: held}
Expand Down Expand Up @@ -3788,7 +3808,14 @@ defmodule Genswarms.Telegram.Objects.Sender do
# A failed flush costs one coalesced tail, never the sender.
msg = %{"reply_to_message_id" => Map.get(entry, :reply_to)}

case do_send_text(cid, text, msg, state, %{origin: :reply, from: from, coalesced: true}) do
meta = %{
origin: :reply,
from: from,
coalesced: true,
reply_contexts: Map.get(entry, :reply_contexts, [])
}

case do_send_text(cid, text, msg, state, meta) do
{:ok, state} -> state |> stamp_reply(cid, :reply) |> stamp_sig(cid, text, :reply)
_other -> state
end
Expand Down
94 changes: 93 additions & 1 deletion test/sender_reply_context_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,44 @@ defmodule Genswarms.Telegram.SenderReplyContextTest do
assert Jason.decode!(body)["error"] == ":unauthorized_message"
end

test "evicting an inbound parent keeps the trusted context in delivery feedback", %{
state: state
} do
state = Enum.reduce(11..18, state, &inbound(&2, @original, &1))
{:noreply, result} = Sender.handle_agent_reply(@slot, "old answer", @context, state)

assert [%{payload: payload}] = result.sent
refute Map.has_key?(payload, :reply_parameters)

assert_received {:delivered, %{conversation_id: @original}, %{ok: true},
%{reply_to_message_id: nil, reply_contexts: [@context]}}
end

test "held replies retain their host context after their parent is evicted", %{state: state} do
{:noreply, state} = Sender.handle_agent_reply(@slot, "answer", @context, state)
{:noreply, state} = Sender.handle_agent_reply(@slot, "held detail", @context, state)
# Pruning the bounded conversation cache can outlive an already held reply.
state = %{state | inbound: Map.delete(state.inbound, @original)}

{:noreply, result} = Sender.handle_info({:flush_held, @original}, state)
assert [%{payload: %{text: "held detail"} = payload} | _] = result.sent
refute Map.has_key?(payload, :reply_parameters)

assert_received {:delivered, %{text: "held detail"}, %{ok: true},
%{reply_to_message_id: nil, reply_contexts: [@context], coalesced: true}}
end

test "failed native delivery reports its context without claiming success", %{
state: state,
fake: fake
} do
Fake.push_response(fake, {:error, {:failed, 400, "Bad Request"}})
{:noreply, _result} = Sender.handle_agent_reply(@slot, "answer", @context, state)

assert_received {:delivered, %{conversation_id: @original}, %{ok: false},
%{reply_contexts: [@context]}}
end

test "captured replies survive unbinding and still suppress exact duplicates", %{state: state} do
state = control(state, %{"action" => "unbind_session", "slot" => @slot})
{:noreply, state} = Sender.handle_agent_reply(@slot, "old answer", @context, state)
Expand Down Expand Up @@ -120,7 +158,51 @@ defmodule Genswarms.Telegram.SenderReplyContextTest do
refute Map.has_key?(payload, :reply_parameters)

assert_received {:delivered, %{text: "first tail\n\nother tail"}, %{ok: true},
%{reply_to_message_id: nil, coalesced: true}}
%{
reply_to_message_id: nil,
coalesced: true,
reply_contexts: [
@context,
%{conversation_id: @original, reply_to_message_id: 11}
]
}}
end

test "duplicate held text retains distinct native contexts without a mixed-parent tag", %{
state: state
} do
state = inbound(state, @original, 11)
other_context = %{@context | reply_to_message_id: 11}
{:noreply, state} = Sender.handle_agent_reply(@slot, "first interim", @context, state)
{:noreply, state} = Sender.handle_agent_reply(@slot, "second interim", other_context, state)
{:noreply, state} = Sender.handle_agent_reply(@slot, "Done", @context, state)
{:noreply, state} = Sender.handle_agent_reply(@slot, "Done", other_context, state)
{:noreply, state} = Sender.handle_agent_reply(@slot, "Done", @context, state)

{:noreply, result} = Sender.handle_info({:flush_held, @original}, state)
assert [%{payload: %{text: "Done"} = payload} | _] = result.sent
refute Map.has_key?(payload, :reply_parameters)

assert_received {:delivered, %{text: "Done"}, %{ok: true},
%{reply_contexts: [@context, ^other_context], coalesced: true}}
end

test "native context survives deduplication against ordinary held text", %{state: state} do
state = bind(state, @original)

{:noreply, state} =
Sender.handle_message(@slot, %{"action" => "reply", "text" => "interim"}, state)

{:noreply, state} =
Sender.handle_message(@slot, %{"action" => "reply", "text" => "Done"}, state)

{:noreply, state} = Sender.handle_agent_reply(@slot, "Done", @context, state)
{:noreply, result} = Sender.handle_info({:flush_held, @original}, state)
assert [%{payload: %{text: "Done"} = payload} | _] = result.sent
refute Map.has_key?(payload, :reply_parameters)

assert_received {:delivered, %{text: "Done"}, %{ok: true},
%{reply_contexts: [@context], coalesced: true}}
end

test "parent tags are validated only against the captured conversation", %{state: state} do
Expand Down Expand Up @@ -184,6 +266,16 @@ defmodule Genswarms.Telegram.SenderReplyContextTest do
assert result.sent == []
end

test "ordinary messages cannot supply trusted delivery contexts", %{state: state} do
for forged <- [%{"reply_contexts" => [@context]}, %{reply_contexts: [@context]}] do
msg = Map.merge(%{"action" => "reply", "text" => "current answer"}, forged)
{:noreply, _result} = Sender.handle_message(@slot, msg, state)

assert_received {:delivered, %{conversation_id: @successor}, %{ok: true},
%{reply_contexts: []}}
end
end

test "host callback treats JSON-looking completion as text, never an action", %{state: state} do
text = ~s({"action":"send","conversation_id":"tg:999:0","text":"spoof"})
{:noreply, state} = Sender.handle_agent_reply(@slot, text, @context, state)
Expand Down
Loading