Skip to content
Merged
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
22 changes: 17 additions & 5 deletions cmd/odek/serve_runs.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@ import (
"fmt"
"net/http"
"sort"
"strconv"
"strings"
"sync"
"sync/atomic"
Expand Down Expand Up @@ -116,9 +117,17 @@ func handleEvents() http.HandlerFunc {
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
return
}
limit := parseIntDefault(r.URL.Query().Get("limit"), 100)
if limit < 1 {
limit = 100
// snapshot() treats limit<=0 as "return everything"; only absent or
// unparseable params get the default. An explicit limit=0 must reach
// snapshot so clients can request the untruncated feed.
limit := 100
if q := r.URL.Query().Get("limit"); q != "" {
if v, err := strconv.Atoi(q); err == nil {
limit = v
}
}
if limit < 0 {
limit = 0
}
if limit > serveEventsCap {
limit = serveEventsCap
Expand Down Expand Up @@ -459,11 +468,14 @@ func (r *serveRun) record(v any) error {
r.Status = "running"
}
case "token", "token_delta":
if c, _ := m["content"].(string); c != "" {
// A still-draining loop must not mutate a terminal run's result —
// the snapshot served by GET /api/runs/{id} would change between
// polls after the run already completed/failed.
if c, _ := m["content"].(string); c != "" && !runStatusTerminal(r.Status) {
r.Result += c
}
case "error":
if msg, _ := m["message"].(string); msg != "" {
if msg, _ := m["message"].(string); msg != "" && !runStatusTerminal(r.Status) {
r.Error = msg
}
case "done":
Expand Down
49 changes: 49 additions & 0 deletions cmd/odek/serve_runs_events_limit_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
package main

import (
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"testing"
"time"

"github.com/BackendStack21/odek/internal/events"
)

// An explicit limit=0 must reach the snapshot untruncated (documented
// "return everything"), and negatives are treated as 0 — previously both
// were silently rewritten to the default 100.
func TestHandleEvents_ExplicitZeroLimitReturnsAll(t *testing.T) {
serveEvents.reset()
t.Cleanup(serveEvents.reset)
now := time.Now().UTC()
for i := 0; i < 3; i++ {
serveEvents.add(events.Event{Type: "e", RunID: fmt.Sprintf("r%d", i), Timestamp: now})
}

decode := func(req *http.Request) (count int) {
w := httptest.NewRecorder()
handleEvents()(w, req)
var body struct {
Count int `json:"count"`
}
if err := json.NewDecoder(w.Body).Decode(&body); err != nil {
t.Fatal(err)
}
return body.Count
}

if got := decode(httptest.NewRequest(http.MethodGet, "/api/events?limit=0", nil)); got != 3 {
t.Errorf("limit=0 returned %d events, want all 3", got)
}
if got := decode(httptest.NewRequest(http.MethodGet, "/api/events?limit=-5", nil)); got != 3 {
t.Errorf("limit=-5 returned %d events, want all 3 (treated as 0)", got)
}
if got := decode(httptest.NewRequest(http.MethodGet, "/api/events?limit=2", nil)); got != 2 {
t.Errorf("limit=2 returned %d events, want 2", got)
}
if got := decode(httptest.NewRequest(http.MethodGet, "/api/events", nil)); got != 3 {
t.Errorf("absent limit returned %d events, want default (all 3 fit)", got)
}
}
9 changes: 8 additions & 1 deletion cmd/odek/wsapprover.go
Original file line number Diff line number Diff line change
Expand Up @@ -182,7 +182,10 @@ func (a *wsApprover) PromptCommand(cls danger.RiskClass, cmd, description string
// the map write below (or SetTrustAll) is a data race that can fatally
// crash serve with "concurrent map read and map write".
a.mu.Lock()
trusted := a.approveAll[cls] || a.trustAll
// Trust-all honors the same class gate as the per-class trust shortcut:
// destructive/blocked/unknown prompts always need a real user decision,
// even after a blanket trust grant from an earlier benign batch.
trusted := a.approveAll[cls] || (a.trustAll && danger.TrustShortcutAllowed(cls))
a.mu.Unlock()
if trusted {
return nil
Expand Down Expand Up @@ -275,6 +278,10 @@ func (a *wsApprover) PromptCommand(cls danger.RiskClass, cmd, description string
a.mu.Lock()
a.approveAll[cls] = true
a.mu.Unlock()
// Trust grants count toward approval friction, same as the TTY
// path — otherwise reflexive trust-tapping never trips the
// fatigue guard on the WebUI.
a.recordApproval(cls)
return nil
default:
return fmt.Errorf("operation denied by user: %s", cmd)
Expand Down
61 changes: 61 additions & 0 deletions cmd/odek/wsapprover_trustall_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,61 @@
package main

import (
"testing"
"time"

"github.com/BackendStack21/odek/internal/danger"
)

// trustAll must not auto-approve classes that never allow trust shortcuts
// (destructive, blocked, unknown, ...). A benign batch approval that grants
// SetTrustAll must still surface a per-call prompt for a destructive command.
func TestWSApprover_TrustAllGatedByTrustShortcutAllowed(t *testing.T) {
a := newWSApprover(func(v any) error { return nil })
a.SetTrustAll(true)

prompted := make(chan struct{}, 1)
go func() {
// Wait for the prompt to be registered, then answer it.
deadline := time.After(3 * time.Second)
for {
a.mu.Lock()
var id string
for k := range a.pending {
id = k
}
a.mu.Unlock()
if id != "" {
a.HandleResponse(id, "approve")
close(prompted)
return
}
select {
case <-deadline:
return
default:
time.Sleep(5 * time.Millisecond)
}
}
}()

err := a.PromptCommand(danger.Destructive, "rm -rf /tmp/x", "test")
if err != nil {
t.Fatalf("PromptCommand returned error: %v", err)
}

select {
case <-prompted:
// Good: a real prompt was shown and answered.
default:
t.Fatal("destructive command was auto-approved by trustAll without prompting")
}

// And the destructive class must not have been silently trusted.
a.mu.Lock()
_, trusted := a.approveAll[danger.Destructive]
a.mu.Unlock()
if trusted {
t.Error("destructive class must never be cached as approved")
}
}
7 changes: 5 additions & 2 deletions docs/WEBUI.md
Original file line number Diff line number Diff line change
Expand Up @@ -527,8 +527,11 @@ classifications — never prompt or completion content.
### `GET /api/events?limit=&run_id=&session_id=`

Recent `odek.event/v1` runtime events (ring of 500, oldest-first, filtered).
Event payloads carry SHA-256 arg hashes and redacted fields only — never raw
tool arguments. Every WS prompt and REST run feeds the same ring.
`limit` defaults to 100; an explicit `limit=0` returns the full filtered
ring (negative values are treated as 0); values above the ring size clamp
to the ring. Event payloads carry SHA-256 arg hashes and redacted fields
only — never raw tool arguments. Every WS prompt and REST run feeds the
same ring.

### `GET /api/usage`

Expand Down
15 changes: 10 additions & 5 deletions internal/fsatomic/fsatomic.go
Original file line number Diff line number Diff line change
Expand Up @@ -66,11 +66,16 @@ func WriteFile(path string, data []byte, perm os.FileMode) (err error) {
tmp = "" // renamed — no longer ours to remove

// Make the rename itself durable by fsyncing the parent directory.
if d, derr := os.Open(dir); derr == nil {
defer d.Close()
if err := d.Sync(); err != nil {
return fmt.Errorf("fsatomic: fsync dir: %w", err)
}
// A failure to open the directory means durability cannot be
// guaranteed for the rename — report it instead of silently
// returning success.
d, derr := os.Open(dir)
if derr != nil {
return fmt.Errorf("fsatomic: open dir: %w", derr)
}
defer d.Close()
if err := d.Sync(); err != nil {
return fmt.Errorf("fsatomic: fsync dir: %w", err)
}
return nil
}
2 changes: 1 addition & 1 deletion internal/resource/resource.go
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,7 @@ func escapeGlob(s string) string {
var b strings.Builder
for _, r := range s {
switch r {
case '*', '?', '[', ']', '^', '{', '}':
case '\\', '*', '?', '[', ']', '^', '{', '}':
b.WriteByte('\\')
}
b.WriteRune(r)
Expand Down
41 changes: 29 additions & 12 deletions internal/telegram/approver.go
Original file line number Diff line number Diff line change
Expand Up @@ -169,9 +169,11 @@ func allowTrustForClass(cls danger.RiskClass) bool {
}

func (a *TelegramApprover) PromptCommand(cls danger.RiskClass, cmd, description string) error {
// Check session trust cache
// Check session trust cache. Trust-all is gated by the same class rule
// as the per-class trust shortcut: destructive/blocked/unknown prompts
// always require an explicit user decision.
a.mu.Lock()
if a.trusted[cls] || a.trustAll {
if a.trusted[cls] || (a.trustAll && danger.TrustShortcutAllowed(cls)) {
a.mu.Unlock()
return nil
}
Expand Down Expand Up @@ -222,20 +224,27 @@ func (a *TelegramApprover) PromptCommand(cls danger.RiskClass, cmd, description
}
markup := InlineKeyboardMarkup{InlineKeyboard: keyboard}

msg, err := a.bot.SendMessage(a.ChatID, text, &SendOpts{
// Register the pending request BEFORE sending so a fast tap racing the
// SendMessage round-trip is not silently dropped (the callback used to
// find no pending entry yet, acknowledge with a toast, and the approval
// then timed out despite the user having acted).
pr := &pendingRequest{resp: make(chan string, 1), userID: a.userID, class: cls, allowTrust: allowTrust}
a.mu.Lock()
a.pending[id] = pr
a.mu.Unlock()

if msg, err := a.bot.SendMessage(a.ChatID, text, &SendOpts{
ParseMode: ParseModeMarkdownV2,
ReplyMarkup: &markup,
})
if err != nil {
}); err != nil {
a.mu.Lock()
delete(a.pending, id)
a.mu.Unlock()
return fmt.Errorf("telegram approver: send prompt: %w", err)
} else {
pr.messageID = msg.ID
}

// Register the pending request with message ID and originating user.
pr := &pendingRequest{resp: make(chan string, 1), messageID: msg.ID, userID: a.userID, class: cls, allowTrust: allowTrust}
a.mu.Lock()
a.pending[id] = pr
a.mu.Unlock()

defer func() {
a.mu.Lock()
delete(a.pending, id)
Expand Down Expand Up @@ -345,7 +354,15 @@ func (a *TelegramApprover) HandleCallback(data string, userID int64) bool {
if pr.userID != 0 && pr.userID != userID {
return true
}
pr.resp <- action
// Non-blocking send: the buffered channel has a single-shot reader.
// A duplicate callback (double-tap, Telegram retry) arriving after
// the reader consumed the action must not block HandleCallback —
// it runs on the serialized update loop, so a blocking send would
// freeze the whole bot for every chat.
select {
case pr.resp <- action:
default:
}
}

return true
Expand Down
41 changes: 41 additions & 0 deletions internal/telegram/approver_deadlock_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
package telegram

import (
"testing"
"time"
)

// A duplicate approval callback (double-tap, Telegram retry) must never
// block HandleCallback: the reader of pr.resp consumes at most one action,
// so a second send on the buffered channel used to block the serialized
// update loop forever, freezing the whole bot.
func TestHandleCallback_DuplicateDoesNotBlock(t *testing.T) {
ts := testServer(t, nil)
defer ts.Close()
bot := testBot(t, ts)

a := NewTelegramApprover(bot, 1, 0)
id := a.newID()
pr := &pendingRequest{resp: make(chan string, 1)}
a.mu.Lock()
a.pending[id] = pr
a.mu.Unlock()

done := make(chan struct{})
go func() {
a.HandleCallback(cbPrefixApprove+id, 0) // first — fills the buffer
a.HandleCallback(cbPrefixApprove+id, 0) // duplicate — must not block
close(done)
}()

select {
case <-done:
case <-time.After(5 * time.Second):
t.Fatal("HandleCallback blocked on duplicate callback (deadlock)")
}

action := <-pr.resp
if action != "approve" {
t.Errorf("action = %q, want %q", action, "approve")
}
}
Loading