No deadline survives the turn path, from mavweb down to llama-server (#188)
Co-authored-by: claude <no-reply@agents.claude.kvmx.ru> Co-committed-by: claude <no-reply@agents.claude.kvmx.ru>
This commit was merged in pull request #188.
This commit is contained in:
@@ -571,7 +571,7 @@ func (h *reactiveHandler) finishClarified(ctx context.Context, dec router.Decisi
|
|||||||
}
|
}
|
||||||
reply := h.applyAction(ctx, dec)
|
reply := h.applyAction(ctx, dec)
|
||||||
if reply == "" {
|
if reply == "" {
|
||||||
reply = h.replier.Reply(dec)
|
reply = h.replier.Reply(ctx, dec)
|
||||||
}
|
}
|
||||||
if reply == "" {
|
if reply == "" {
|
||||||
// Belt: an empty reply here would be a silent drop.
|
// Belt: an empty reply here would be a silent drop.
|
||||||
|
|||||||
@@ -22,7 +22,7 @@ func newLLMReplier(c phraser.Completer, block func() string) *llmReplier {
|
|||||||
|
|
||||||
// Reply never fails: a clarify, a model error and an unusable generation all
|
// Reply never fails: a clarify, a model error and an unusable generation all
|
||||||
// answer from the stub, which is what keeps a turn from breaking on the model.
|
// answer from the stub, which is what keeps a turn from breaking on the model.
|
||||||
func (r *llmReplier) Reply(d router.Decision) string {
|
func (r *llmReplier) Reply(ctx context.Context, d router.Decision) string {
|
||||||
if d.Clarify {
|
if d.Clarify {
|
||||||
// The deck, not the stub's single sentence: a clarify she cannot turn
|
// The deck, not the stub's single sentence: a clarify she cannot turn
|
||||||
// into a question is the line he hears most often when she misses him,
|
// into a question is the line he hears most often when she misses him,
|
||||||
@@ -39,14 +39,14 @@ func (r *llmReplier) Reply(d router.Decision) string {
|
|||||||
// что ты выпел стакан воды" for "я выпил воды".
|
// что ты выпел стакан воды" for "я выпил воды".
|
||||||
return phraser.FactAck(d.Utterance)
|
return phraser.FactAck(d.Utterance)
|
||||||
}
|
}
|
||||||
out, err := r.p.PhraseReply(context.Background(), d)
|
out, err := r.p.PhraseReply(ctx, d)
|
||||||
if err != nil || out == "" {
|
if err != nil || out == "" {
|
||||||
return r.stub.Reply(d)
|
return r.stub.Reply(ctx, d)
|
||||||
}
|
}
|
||||||
// The persona checks, on the live path (personaguard.go). A reply that
|
// The persona checks, on the live path (personaguard.go). A reply that
|
||||||
// leaks reasoning or calls him "вы" is worse than a flat one.
|
// leaks reasoning or calls him "вы" is worse than a flat one.
|
||||||
if _, ok := guardSpoken("reply", out); !ok {
|
if _, ok := guardSpoken("reply", out); !ok {
|
||||||
return r.stub.Reply(d)
|
return r.stub.Reply(ctx, d)
|
||||||
}
|
}
|
||||||
return out
|
return out
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -22,7 +22,7 @@ func (s stubCompleter) Complete(_ context.Context, _ llm.Req) (string, error) {
|
|||||||
|
|
||||||
func TestLLMReplierPassesTheModelReplyThrough(t *testing.T) {
|
func TestLLMReplierPassesTheModelReplyThrough(t *testing.T) {
|
||||||
r := newLLMReplier(stubCompleter{out: `{"response":"записала, кофе закончился","mood":"neutral"}`}, nil)
|
r := newLLMReplier(stubCompleter{out: `{"response":"записала, кофе закончился","mood":"neutral"}`}, nil)
|
||||||
got := r.Reply(router.Decision{Intent: router.IntentNote, Slots: router.Slots{Text: "кофе закончился"}})
|
got := r.Reply(context.Background(), router.Decision{Intent: router.IntentNote, Slots: router.Slots{Text: "кофе закончился"}})
|
||||||
if got != "записала, кофе закончился" {
|
if got != "записала, кофе закончился" {
|
||||||
t.Errorf("got %q, want %q", got, "записала, кофе закончился")
|
t.Errorf("got %q, want %q", got, "записала, кофе закончился")
|
||||||
}
|
}
|
||||||
@@ -42,7 +42,7 @@ func TestLLMReplierFallsBackToStubOnEmpty(t *testing.T) {
|
|||||||
// the clarify deck rather than the stub's single sentence.
|
// the clarify deck rather than the stub's single sentence.
|
||||||
func TestLLMReplierClarifyReadsTheDeck(t *testing.T) {
|
func TestLLMReplierClarifyReadsTheDeck(t *testing.T) {
|
||||||
r := newLLMReplier(stubCompleter{out: "я всё поняла"}, nil)
|
r := newLLMReplier(stubCompleter{out: "я всё поняла"}, nil)
|
||||||
got := r.Reply(router.Decision{Clarify: true, Utterance: "мгм"})
|
got := r.Reply(context.Background(), router.Decision{Clarify: true, Utterance: "мгм"})
|
||||||
if got == "я всё поняла" {
|
if got == "я всё поняла" {
|
||||||
t.Fatal("a clarify must not be phrased by the model")
|
t.Fatal("a clarify must not be phrased by the model")
|
||||||
}
|
}
|
||||||
@@ -50,7 +50,7 @@ func TestLLMReplierClarifyReadsTheDeck(t *testing.T) {
|
|||||||
t.Errorf("on clarify: got %q, want %q", got, want)
|
t.Errorf("on clarify: got %q, want %q", got, want)
|
||||||
}
|
}
|
||||||
// Two different misses do not sound identical.
|
// Two different misses do not sound identical.
|
||||||
if same := r.Reply(router.Decision{Clarify: true, Utterance: "а"}); same == got {
|
if same := r.Reply(context.Background(), router.Decision{Clarify: true, Utterance: "а"}); same == got {
|
||||||
t.Log("two utterances hashed to the same line, which is allowed but should be rare")
|
t.Log("two utterances hashed to the same line, which is allowed but should be rare")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -60,14 +60,14 @@ func TestLLMReplierClarifyReadsTheDeck(t *testing.T) {
|
|||||||
// produce, which is the same claim without pinning one wording.
|
// produce, which is the same claim without pinning one wording.
|
||||||
func assertAck(t *testing.T, r *llmReplier, d router.Decision, key, what string) {
|
func assertAck(t *testing.T, r *llmReplier, d router.Decision, key, what string) {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
if got := r.Reply(d); !phraser.IsAck(key, nil, got) {
|
if got := r.Reply(context.Background(), d); !phraser.IsAck(key, nil, got) {
|
||||||
t.Errorf("on %s: got %q, want a %q line", what, got, key)
|
t.Errorf("on %s: got %q, want a %q line", what, got, key)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func assertStub(t *testing.T, r *llmReplier, d router.Decision, what string) {
|
func assertStub(t *testing.T, r *llmReplier, d router.Decision, what string) {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
got, want := r.Reply(d), voice.NewStubReplier().Reply(d)
|
got, want := r.Reply(context.Background(), d), voice.NewStubReplier().Reply(context.Background(), d)
|
||||||
if got != want {
|
if got != want {
|
||||||
t.Errorf("on %s: got %q, want stub %q", what, got, want)
|
t.Errorf("on %s: got %q, want stub %q", what, got, want)
|
||||||
}
|
}
|
||||||
|
|||||||
+1
-1
@@ -458,7 +458,7 @@ func (h *reactiveHandler) runTurn(ctx context.Context, text string, src turnSour
|
|||||||
|
|
||||||
// 9. replier — phrase the reply across the router decision.
|
// 9. replier — phrase the reply across the router decision.
|
||||||
if replyText == "" {
|
if replyText == "" {
|
||||||
replyText = h.replier.Reply(dec)
|
replyText = h.replier.Reply(ctx, dec)
|
||||||
}
|
}
|
||||||
return withNotice(expiredNotice, replyText)
|
return withNotice(expiredNotice, replyText)
|
||||||
}
|
}
|
||||||
|
|||||||
+19
-1
@@ -58,6 +58,12 @@ func main() {
|
|||||||
// mutex, so sharing the connection would freeze every other page for the
|
// mutex, so sharing the connection would freeze every other page for the
|
||||||
// length of the load. See handleModels.
|
// length of the load. See handleModels.
|
||||||
var swapConn modelController
|
var swapConn modelController
|
||||||
|
// turnConn — a third connection, for POST /api/chat and nothing else, for
|
||||||
|
// the same reason /models has one (V-638). A chat turn routes, phrases and
|
||||||
|
// may act, bounded only by phraser.timeout at 60s, and every other handler
|
||||||
|
// on this server queues behind it on the shared client's one mutex. Nil ⇒
|
||||||
|
// chat shares the main connection, which is how it behaved before.
|
||||||
|
var turnConn ipc.CoreAPI
|
||||||
if *coreSock != "" {
|
if *coreSock != "" {
|
||||||
c, err := ipc.DialWait(*coreSock, 60*time.Second)
|
c, err := ipc.DialWait(*coreSock, 60*time.Second)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -71,6 +77,12 @@ func main() {
|
|||||||
defer sc.Close()
|
defer sc.Close()
|
||||||
swapConn = sc
|
swapConn = sc
|
||||||
}
|
}
|
||||||
|
if tc, err := ipc.Dial(*coreSock); err != nil {
|
||||||
|
log.Printf("chat: third core connection failed (%v) — /api/chat will share the main one and a turn will block the other pages", err)
|
||||||
|
} else {
|
||||||
|
defer tc.Close()
|
||||||
|
turnConn = tc
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// stepUpSession stays nil unless the passkey endpoints are wired below — it
|
// stepUpSession stays nil unless the passkey endpoints are wired below — it
|
||||||
@@ -208,7 +220,13 @@ func main() {
|
|||||||
// decides how every utterance is routed and how every reply is worded.
|
// decides how every utterance is routed and how every reply is worded.
|
||||||
mux.HandleFunc("/tools", gatedPage(handleTools))
|
mux.HandleFunc("/tools", gatedPage(handleTools))
|
||||||
mux.HandleFunc("/routines", gatedPage(handleRoutines))
|
mux.HandleFunc("/routines", gatedPage(handleRoutines))
|
||||||
mux.HandleFunc("/api/chat", gatedPage(handleChatAPI))
|
mux.HandleFunc("/api/chat", func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
c := turnConn
|
||||||
|
if c == nil {
|
||||||
|
c = core
|
||||||
|
}
|
||||||
|
handleChatAPI(w, r, c, stepUpSession, *requireStepUp)
|
||||||
|
})
|
||||||
mux.HandleFunc("/api/revert", gatedPage(handleRevert))
|
mux.HandleFunc("/api/revert", gatedPage(handleRevert))
|
||||||
mux.HandleFunc("/api/correct", gatedPage(handleCorrectAPI))
|
mux.HandleFunc("/api/correct", gatedPage(handleCorrectAPI))
|
||||||
mux.HandleFunc("/models", func(w http.ResponseWriter, r *http.Request) {
|
mux.HandleFunc("/models", func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
|||||||
@@ -1,6 +1,14 @@
|
|||||||
# No deadline on the turn path
|
# No deadline on the turn path
|
||||||
|
|
||||||
Last verified: 06-08-2026 @ 06c1cf2
|
Last verified: 06-08-2026 @ 60e64dd
|
||||||
|
|
||||||
|
**All four steps landed on 06-08-2026.** What follows describes the defect as it was and
|
||||||
|
the work as it was planned. Two things came out differently. `Client.Close` read the conn
|
||||||
|
field with no lock while `roundtrip` re-dialed and dropped it. `-race` caught that on the
|
||||||
|
new cancellation test. So the conn field now has a mutex of its own, held only across a
|
||||||
|
read or an assignment. And `/api/ptt` needed nothing: it proxies to the voice port and never
|
||||||
|
touches the shared client, so only `/api/chat` got the extra connection. The pool inside
|
||||||
|
`ipc.Client` is still unbuilt and still waiting on a second module measured queueing.
|
||||||
|
|
||||||
V-638. Sibling of V-607, which is the same class of bug in `internal/worker`.
|
V-638. Sibling of V-607, which is the same class of bug in `internal/worker`.
|
||||||
Reads with `docs/offload.md` and `docs/protocol.md`.
|
Reads with `docs/offload.md` and `docs/protocol.md`.
|
||||||
|
|||||||
@@ -0,0 +1,192 @@
|
|||||||
|
package ipc
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"net"
|
||||||
|
"path/filepath"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
// A cancelled context has to abort a call that is already in flight. It did not
|
||||||
|
// until V-638: call checked ctx once before sending and then blocked in
|
||||||
|
// roundtrip with no connection deadline, so a daemon that read the frame and
|
||||||
|
// never answered parked the caller for as long as the socket stayed open.
|
||||||
|
//
|
||||||
|
// The server here is that daemon: it accepts, reads nothing, replies nothing.
|
||||||
|
|
||||||
|
func deafServer(t *testing.T) string {
|
||||||
|
t.Helper()
|
||||||
|
sock := filepath.Join(t.TempDir(), "deaf.sock")
|
||||||
|
ln, err := net.Listen("unix", sock)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("listen: %v", err)
|
||||||
|
}
|
||||||
|
t.Cleanup(func() { _ = ln.Close() })
|
||||||
|
go func() {
|
||||||
|
for {
|
||||||
|
conn, err := ln.Accept()
|
||||||
|
if err != nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
// Hold it open and say nothing. Closed by the listener cleanup.
|
||||||
|
t.Cleanup(func() { _ = conn.Close() })
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
return sock
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestClientCancelAbortsAReadInFlight(t *testing.T) {
|
||||||
|
c, err := Dial(deafServer(t))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("dial: %v", err)
|
||||||
|
}
|
||||||
|
defer c.Close()
|
||||||
|
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
go func() {
|
||||||
|
time.Sleep(50 * time.Millisecond)
|
||||||
|
cancel()
|
||||||
|
}()
|
||||||
|
|
||||||
|
done := make(chan error, 1)
|
||||||
|
go func() {
|
||||||
|
_, err := c.Ping(ctx)
|
||||||
|
done <- err
|
||||||
|
}()
|
||||||
|
|
||||||
|
select {
|
||||||
|
case err := <-done:
|
||||||
|
// Ping is read-only, so the cancellation is reported as itself rather
|
||||||
|
// than as an ambiguous mutation.
|
||||||
|
if !errors.Is(err, context.Canceled) {
|
||||||
|
t.Errorf("got %v, want context.Canceled", err)
|
||||||
|
}
|
||||||
|
case <-time.After(5 * time.Second):
|
||||||
|
t.Fatal("a cancelled Ping did not return")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A mutation cancelled while awaiting the reply may already have committed, so
|
||||||
|
// it is ErrAmbiguousOutcome and never a retry. That split is the invariant
|
||||||
|
// internal/ipc/maperr_test.go's neighbours rest on.
|
||||||
|
func TestClientCancelLeavesAMutationAmbiguous(t *testing.T) {
|
||||||
|
c, err := Dial(deafServer(t))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("dial: %v", err)
|
||||||
|
}
|
||||||
|
defer c.Close()
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
done := make(chan error, 1)
|
||||||
|
go func() {
|
||||||
|
_, err := c.WriteFact(ctx, WriteFactReq{Key: "water", Value: "drank"})
|
||||||
|
done <- err
|
||||||
|
}()
|
||||||
|
|
||||||
|
select {
|
||||||
|
case err := <-done:
|
||||||
|
if !errors.Is(err, ErrAmbiguousOutcome) {
|
||||||
|
t.Errorf("got %v, want ErrAmbiguousOutcome", err)
|
||||||
|
}
|
||||||
|
case <-time.After(5 * time.Second):
|
||||||
|
t.Fatal("a cancelled WriteFact did not return")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// The deadline itself, with no cancellation: a call on a context with no
|
||||||
|
// deadline used to have no bound at all. This one has one and must respect it.
|
||||||
|
func TestClientDeadlineBoundsACall(t *testing.T) {
|
||||||
|
c, err := Dial(deafServer(t))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("dial: %v", err)
|
||||||
|
}
|
||||||
|
defer c.Close()
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
start := time.Now()
|
||||||
|
if _, err := c.Ping(ctx); err == nil {
|
||||||
|
t.Fatal("a deaf server answered a Ping")
|
||||||
|
}
|
||||||
|
if elapsed := time.Since(start); elapsed > 3*time.Second {
|
||||||
|
t.Errorf("Ping took %v, want the context deadline to bound it", elapsed)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// blockingAPI parks Presence until its context is cancelled and records what
|
||||||
|
// cancelled it. Every other method is the unimplemented floor.
|
||||||
|
type blockingAPI struct {
|
||||||
|
UnimplementedCoreAPI
|
||||||
|
entered chan struct{}
|
||||||
|
err chan error
|
||||||
|
}
|
||||||
|
|
||||||
|
func (b *blockingAPI) Presence(ctx context.Context) (Presence, error) {
|
||||||
|
close(b.entered)
|
||||||
|
<-ctx.Done()
|
||||||
|
b.err <- ctx.Err()
|
||||||
|
return Presence{}, ctx.Err()
|
||||||
|
}
|
||||||
|
|
||||||
|
// serveConn dispatched under context.Background() until V-638, so Close could
|
||||||
|
// only abandon a dispatch in flight and never tell it to stop.
|
||||||
|
func TestServerCloseCancelsADispatchInFlight(t *testing.T) {
|
||||||
|
api := &blockingAPI{entered: make(chan struct{}), err: make(chan error, 1)}
|
||||||
|
srv, err := Listen(filepath.Join(t.TempDir(), "core.sock"), api)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("listen: %v", err)
|
||||||
|
}
|
||||||
|
served := make(chan struct{})
|
||||||
|
go func() { _ = srv.Serve(); close(served) }()
|
||||||
|
|
||||||
|
cli, err := Dial(srv.Path())
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("dial: %v", err)
|
||||||
|
}
|
||||||
|
defer cli.Close()
|
||||||
|
go func() { _, _ = cli.Presence(context.Background()) }()
|
||||||
|
|
||||||
|
select {
|
||||||
|
case <-api.entered:
|
||||||
|
case <-time.After(5 * time.Second):
|
||||||
|
t.Fatal("the handler was never dispatched")
|
||||||
|
}
|
||||||
|
|
||||||
|
_ = srv.Close()
|
||||||
|
<-served
|
||||||
|
select {
|
||||||
|
case got := <-api.err:
|
||||||
|
if !errors.Is(got, context.Canceled) {
|
||||||
|
t.Errorf("handler saw %v, want context.Canceled", got)
|
||||||
|
}
|
||||||
|
case <-time.After(5 * time.Second):
|
||||||
|
t.Fatal("Close did not cancel the dispatch")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// The watchdog closes the conn, and it races the end of the call: a
|
||||||
|
// cancellation landing as the reply arrives can close a conn the call was
|
||||||
|
// already done with. That is survivable either way, because a write to a closed
|
||||||
|
// socket is errWriteLost and errWriteLost re-dials and retries, so this test
|
||||||
|
// passes with or without the drop in roundtrip's defer. What it pins is that
|
||||||
|
// the recovery is real and costs one round trip at most, never an error the
|
||||||
|
// caller sees.
|
||||||
|
func TestClientSurvivesACancelledCall(t *testing.T) {
|
||||||
|
_, _, cli, _ := newServerWithStore(t)
|
||||||
|
|
||||||
|
for i := 0; i < 20; i++ {
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
go cancel() // races the reply on purpose
|
||||||
|
_, _ = cli.Ping(ctx)
|
||||||
|
cancel()
|
||||||
|
|
||||||
|
if _, err := cli.Ping(context.Background()); err != nil {
|
||||||
|
t.Fatalf("call %d after a cancelled one: %v", i, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
+94
-12
@@ -26,9 +26,20 @@ type Client struct {
|
|||||||
conn net.Conn
|
conn net.Conn
|
||||||
path string // the address as configured, kept for errors and logs
|
path string // the address as configured, kept for errors and logs
|
||||||
addr netaddr.Addr // parsed, so a dropped conn can be re-dialed (core restart)
|
addr netaddr.Addr // parsed, so a dropped conn can be re-dialed (core restart)
|
||||||
mu sync.Mutex
|
mu sync.Mutex // one request at a time, so a frame and its reply pair up
|
||||||
|
|
||||||
|
// connMu guards the conn field alone, and is held only across an assignment
|
||||||
|
// or a read. It exists so Close and the cancellation watchdog can reach the
|
||||||
|
// connection without waiting for the call that is holding c.mu (V-638).
|
||||||
|
connMu sync.Mutex
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// defaultCallTimeout bounds a call whose context carries no deadline. It is
|
||||||
|
// the same 120s internal/voice/client.go settles on: long enough for a model
|
||||||
|
// call on a cold resident model, short enough that a daemon which stopped
|
||||||
|
// answering does not park the caller forever.
|
||||||
|
const defaultCallTimeout = 120 * time.Second
|
||||||
|
|
||||||
// errWriteLost marks a conn drop while sending the request frame: the request
|
// errWriteLost marks a conn drop while sending the request frame: the request
|
||||||
// never reached the server (or the server never saw a complete frame), so
|
// never reached the server (or the server never saw a complete frame), so
|
||||||
// retrying is always safe regardless of method — nothing was applied to
|
// retrying is always safe regardless of method — nothing was applied to
|
||||||
@@ -103,11 +114,22 @@ func Dial(path string) (*Client, error) {
|
|||||||
return &Client{conn: c, path: path, addr: addr}, nil
|
return &Client{conn: c, path: path, addr: addr}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Close closes the connection out from under a call in flight, on purpose: a
|
||||||
|
// shutdown must not wait out a parked read. It takes connMu and never c.mu, so
|
||||||
|
// it cannot block behind the call it is interrupting.
|
||||||
|
//
|
||||||
|
// The lock is taken and released by hand, around the two field accesses and
|
||||||
|
// nothing else. The socket close happens outside it, because a close on a tcp
|
||||||
|
// conn can block and connMu is on the path of every call.
|
||||||
func (c *Client) Close() error {
|
func (c *Client) Close() error {
|
||||||
if c.conn == nil {
|
c.connMu.Lock()
|
||||||
|
conn := c.conn
|
||||||
|
c.conn = nil
|
||||||
|
c.connMu.Unlock()
|
||||||
|
if conn == nil {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
return c.conn.Close()
|
return conn.Close()
|
||||||
}
|
}
|
||||||
|
|
||||||
// DialWait is Dial with patience: it retries with capped backoff until the
|
// DialWait is Dial with patience: it retries with capped backoff until the
|
||||||
@@ -163,16 +185,26 @@ func (c *Client) call(ctx context.Context, m Method, params, result any) error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
var resp Response
|
var resp Response
|
||||||
err := c.roundtrip(m, raw, &resp)
|
err := c.roundtrip(ctx, m, raw, &resp)
|
||||||
switch {
|
switch {
|
||||||
case errors.Is(err, errWriteLost):
|
case errors.Is(err, errWriteLost):
|
||||||
// The request never left; a duplicate send can't double-apply.
|
// The request never left; a duplicate send can't double-apply.
|
||||||
// Redial (roundtrip re-dials on a nil conn) and retry exactly once.
|
// Redial (roundtrip re-dials on a nil conn) and retry exactly once.
|
||||||
err = c.roundtrip(m, raw, &resp)
|
// Not when the caller has given up — a retry would only be a second
|
||||||
|
// frame nobody is waiting for.
|
||||||
|
if ctx.Err() == nil {
|
||||||
|
err = c.roundtrip(ctx, m, raw, &resp)
|
||||||
|
}
|
||||||
case errors.Is(err, errReadLost):
|
case errors.Is(err, errReadLost):
|
||||||
if readOnlyMethods[m] {
|
if readOnlyMethods[m] {
|
||||||
|
if ctx.Err() != nil {
|
||||||
|
// The caller cancelled the read it was waiting for. Nothing
|
||||||
|
// was applied, so this is the cancellation and not an
|
||||||
|
// ambiguity.
|
||||||
|
return ctx.Err()
|
||||||
|
}
|
||||||
// A duplicate read can't double-apply either — safe to replay.
|
// A duplicate read can't double-apply either — safe to replay.
|
||||||
err = c.roundtrip(m, raw, &resp)
|
err = c.roundtrip(ctx, m, raw, &resp)
|
||||||
} else {
|
} else {
|
||||||
// The mutation may have already committed server-side. Do not
|
// The mutation may have already committed server-side. Do not
|
||||||
// retry: report the ambiguity instead of guessing.
|
// retry: report the ambiguity instead of guessing.
|
||||||
@@ -199,19 +231,55 @@ func (c *Client) call(ctx context.Context, m Method, params, result any) error {
|
|||||||
// failure is wrapped in errReadLost (ambiguous — call() only retries it for
|
// failure is wrapped in errReadLost (ambiguous — call() only retries it for
|
||||||
// read-only methods). Either way a failed conn is dropped so the next call
|
// read-only methods). Either way a failed conn is dropped so the next call
|
||||||
// re-dials clean. Caller holds c.mu.
|
// re-dials clean. Caller holds c.mu.
|
||||||
func (c *Client) roundtrip(m Method, raw json.RawMessage, resp *Response) error {
|
//
|
||||||
if c.conn == nil {
|
// The connection carries a deadline derived from ctx, falling back to
|
||||||
conn, err := netaddr.Dial(c.addr)
|
// defaultCallTimeout, and a watchdog closes it if ctx is cancelled mid-call
|
||||||
|
// (V-638). Before that a daemon which stopped answering parked the caller for
|
||||||
|
// as long as the socket stayed open. The watchdog closes the conn rather than
|
||||||
|
// calling drop, because drop wants c.mu and the caller is holding it — the
|
||||||
|
// closed socket fails the read, and roundtrip drops it on the way out.
|
||||||
|
func (c *Client) roundtrip(ctx context.Context, m Method, raw json.RawMessage, resp *Response) error {
|
||||||
|
conn := c.currentConn()
|
||||||
|
if conn == nil {
|
||||||
|
dialed, err := netaddr.Dial(c.addr)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("%w: dial %s: %v", errWriteLost, c.addr, err)
|
return fmt.Errorf("%w: dial %s: %v", errWriteLost, c.addr, err)
|
||||||
}
|
}
|
||||||
c.conn = conn
|
c.setConn(dialed)
|
||||||
|
conn = dialed
|
||||||
}
|
}
|
||||||
if err := writeFrame(c.conn, Request{Method: m, Params: raw}); err != nil {
|
if dl, ok := ctx.Deadline(); ok {
|
||||||
|
_ = conn.SetDeadline(dl)
|
||||||
|
} else {
|
||||||
|
_ = conn.SetDeadline(time.Now().Add(defaultCallTimeout))
|
||||||
|
}
|
||||||
|
defer conn.SetDeadline(time.Time{})
|
||||||
|
|
||||||
|
// The watchdog and the end of the call race by construction: a cancellation
|
||||||
|
// landing just as the reply arrives can close a conn this call is already
|
||||||
|
// done with, and c.conn would still point at the closed socket. So a call
|
||||||
|
// whose context ended does not leave the conn behind for the next one,
|
||||||
|
// whichever of the two got there first.
|
||||||
|
done := make(chan struct{})
|
||||||
|
defer func() {
|
||||||
|
close(done)
|
||||||
|
if ctx.Err() != nil {
|
||||||
|
c.drop()
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
go func() {
|
||||||
|
select {
|
||||||
|
case <-ctx.Done():
|
||||||
|
_ = conn.Close()
|
||||||
|
case <-done:
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
if err := writeFrame(conn, Request{Method: m, Params: raw}); err != nil {
|
||||||
c.drop()
|
c.drop()
|
||||||
return fmt.Errorf("%w: %v", errWriteLost, err)
|
return fmt.Errorf("%w: %v", errWriteLost, err)
|
||||||
}
|
}
|
||||||
if err := readFrame(c.conn, resp); err != nil {
|
if err := readFrame(conn, resp); err != nil {
|
||||||
c.drop()
|
c.drop()
|
||||||
return fmt.Errorf("%w: %v", errReadLost, err)
|
return fmt.Errorf("%w: %v", errReadLost, err)
|
||||||
}
|
}
|
||||||
@@ -220,12 +288,26 @@ func (c *Client) roundtrip(m Method, raw json.RawMessage, resp *Response) error
|
|||||||
|
|
||||||
// drop closes and forgets the current conn so the next call re-dials.
|
// drop closes and forgets the current conn so the next call re-dials.
|
||||||
func (c *Client) drop() {
|
func (c *Client) drop() {
|
||||||
|
c.connMu.Lock()
|
||||||
|
defer c.connMu.Unlock()
|
||||||
if c.conn != nil {
|
if c.conn != nil {
|
||||||
_ = c.conn.Close()
|
_ = c.conn.Close()
|
||||||
c.conn = nil
|
c.conn = nil
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (c *Client) currentConn() net.Conn {
|
||||||
|
c.connMu.Lock()
|
||||||
|
defer c.connMu.Unlock()
|
||||||
|
return c.conn
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *Client) setConn(conn net.Conn) {
|
||||||
|
c.connMu.Lock()
|
||||||
|
defer c.connMu.Unlock()
|
||||||
|
c.conn = conn
|
||||||
|
}
|
||||||
|
|
||||||
// hydrate rehydrates a wire RpcError into the matching package sentinel. The
|
// hydrate rehydrates a wire RpcError into the matching package sentinel. The
|
||||||
// code↔sentinel table is the only place the wire "knows" about errors; keep it
|
// code↔sentinel table is the only place the wire "knows" about errors; keep it
|
||||||
// in sync with codeOf in wire.go.
|
// in sync with codeOf in wire.go.
|
||||||
|
|||||||
+31
-5
@@ -30,6 +30,14 @@ type Server struct {
|
|||||||
done chan struct{}
|
done chan struct{}
|
||||||
accept sync.Mutex // guards wg.Add vs Close's wg.Wait sequence
|
accept sync.Mutex // guards wg.Add vs Close's wg.Wait sequence
|
||||||
|
|
||||||
|
// ctx — server-scoped, cancelled by Close, and the parent of every request
|
||||||
|
// context. serveConn dispatched under context.Background() until V-638, so
|
||||||
|
// a dispatch in flight during shutdown could not be told to stop and the
|
||||||
|
// closeGrace below could only abandon it. Cancelling gives a handler that
|
||||||
|
// respects its context the chance to return instead.
|
||||||
|
ctx context.Context
|
||||||
|
cancel context.CancelFunc
|
||||||
|
|
||||||
// conns — every accepted connection still being served. Close needs these
|
// conns — every accepted connection still being served. Close needs these
|
||||||
// because closing the listener does nothing to a connection already
|
// because closing the listener does nothing to a connection already
|
||||||
// accepted: serveConn is parked in readFrame waiting for a peer that may
|
// accepted: serveConn is parked in readFrame waiting for a peer that may
|
||||||
@@ -208,11 +216,14 @@ func Listen(path string, api CoreAPI) (*Server, error) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
s := &Server{
|
s := &Server{
|
||||||
path: path,
|
path: path,
|
||||||
addr: addr,
|
addr: addr,
|
||||||
ln: ln,
|
ln: ln,
|
||||||
done: make(chan struct{}),
|
done: make(chan struct{}),
|
||||||
|
ctx: ctx,
|
||||||
|
cancel: cancel,
|
||||||
}
|
}
|
||||||
s.api.Store(api)
|
s.api.Store(api)
|
||||||
return s, nil
|
return s, nil
|
||||||
@@ -250,7 +261,10 @@ func (s *Server) Serve() error {
|
|||||||
|
|
||||||
func (s *Server) serveConn(c net.Conn) {
|
func (s *Server) serveConn(c net.Conn) {
|
||||||
caller, callerOK := peerCaller(c)
|
caller, callerOK := peerCaller(c)
|
||||||
ctx := context.Background()
|
// Derived from the server's, so Close cancels a dispatch in flight, and
|
||||||
|
// cancelled when this conn ends so nothing a handler spawned outlives it.
|
||||||
|
ctx, cancel := context.WithCancel(s.serverContext())
|
||||||
|
defer cancel()
|
||||||
if callerOK {
|
if callerOK {
|
||||||
ctx = WithCaller(ctx, caller)
|
ctx = WithCaller(ctx, caller)
|
||||||
}
|
}
|
||||||
@@ -274,6 +288,15 @@ func (s *Server) serveConn(c net.Conn) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// serverContext is s.ctx, or Background for a Server built as a zero value
|
||||||
|
// rather than by Listen (the wiring tests do that).
|
||||||
|
func (s *Server) serverContext() context.Context {
|
||||||
|
if s.ctx == nil {
|
||||||
|
return context.Background()
|
||||||
|
}
|
||||||
|
return s.ctx
|
||||||
|
}
|
||||||
|
|
||||||
func (s *Server) safeDispatch(ctx context.Context, req Request) (result json.RawMessage, err error) {
|
func (s *Server) safeDispatch(ctx context.Context, req Request) (result json.RawMessage, err error) {
|
||||||
defer func() {
|
defer func() {
|
||||||
if r := recover(); r != nil {
|
if r := recover(); r != nil {
|
||||||
@@ -720,6 +743,9 @@ func (s *Server) Close() error {
|
|||||||
default:
|
default:
|
||||||
close(s.done)
|
close(s.done)
|
||||||
}
|
}
|
||||||
|
if s.cancel != nil {
|
||||||
|
s.cancel()
|
||||||
|
}
|
||||||
err := s.ln.Close()
|
err := s.ln.Close()
|
||||||
// Closing the listener stops new connections; it does nothing to the ones
|
// Closing the listener stops new connections; it does nothing to the ones
|
||||||
// already accepted. Close those too, or every serveConn parked in readFrame
|
// already accepted. Close those too, or every serveConn parked in readFrame
|
||||||
|
|||||||
@@ -26,6 +26,8 @@
|
|||||||
package voice
|
package voice
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
|
|
||||||
"github.com/kami/maven/internal/phraser"
|
"github.com/kami/maven/internal/phraser"
|
||||||
"github.com/kami/maven/internal/router"
|
"github.com/kami/maven/internal/router"
|
||||||
)
|
)
|
||||||
@@ -38,8 +40,10 @@ import (
|
|||||||
// decision's Intent + Slots + Clarify. The Intent largely names the reply
|
// decision's Intent + Slots + Clarify. The Intent largely names the reply
|
||||||
// shape (act/reminder/fact/note/query/clarify); the Slots carry the
|
// shape (act/reminder/fact/note/query/clarify); the Slots carry the
|
||||||
// specifics that personalise it ("got it: water at 14:00").
|
// specifics that personalise it ("got it: water at 14:00").
|
||||||
|
// The context is the turn's, and it is the only bound an LLM-backed impl has
|
||||||
|
// besides the phraser timeout (V-638). A floor impl ignores it.
|
||||||
type Replier interface {
|
type Replier interface {
|
||||||
Reply(d router.Decision) string
|
Reply(ctx context.Context, d router.Decision) string
|
||||||
}
|
}
|
||||||
|
|
||||||
// StubReplier — the deterministic, no-model floor. Canned per intent;
|
// StubReplier — the deterministic, no-model floor. Canned per intent;
|
||||||
@@ -54,7 +58,8 @@ func NewStubReplier() *StubReplier { return &StubReplier{} }
|
|||||||
|
|
||||||
// Reply dispatches on Intent + Clarify. Each branch is short; the LLM impl
|
// Reply dispatches on Intent + Clarify. Each branch is short; the LLM impl
|
||||||
// will replace this with prompted text and the same dispatch shape.
|
// will replace this with prompted text and the same dispatch shape.
|
||||||
func (s *StubReplier) Reply(d router.Decision) string {
|
// It makes no model call, so the context is unused.
|
||||||
|
func (s *StubReplier) Reply(_ context.Context, d router.Decision) string {
|
||||||
if d.Clarify {
|
if d.Clarify {
|
||||||
return "не совсем поняла — можешь переформулировать?"
|
return "не совсем поняла — можешь переформулировать?"
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user