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
5 changes: 5 additions & 0 deletions edge-server/internal/api/handlers.go
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,11 @@ type Handler struct {
PermissionRegistry *permission.PermissionRegistry
PermissionBroker *adapters.PermissionDecisionBroker

// permissionReceipts provides bounded warm-replay receipts for modern
// POST /v1/permissions/decide controls. Lazy and nil-safe in production;
// it is process-local and never authoritative.
permissionReceipts *permissionReceiptCache

// PlanApprovalBroker manages pending orchestrator plans and connects
// them to user approval/rejection decisions (P0 #3: plan confirmation gate).
PlanApprovalBroker *orchestrator.PlanApprovalBroker
Expand Down
143 changes: 109 additions & 34 deletions edge-server/internal/api/handlers_approvals.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,26 +10,35 @@ import (
"github.com/agenthub/edge-server/internal/permission"
)

// Handler holds dependencies for HTTP and WebSocket handlers.
// permissionDecideRequest is the POST /v1/permissions/decide payload. controlId
// and hubTaskId are optional together: when either is present both are required,
// and the request is modern and idempotent; when both are absent the existing
// one-shot receiver semantics are unchanged.
type permissionDecideRequest struct {
ControlID string `json:"controlId"`
HubTaskID string `json:"hubTaskId"`
RunID string `json:"runId"`
RequestID string `json:"requestId"`
Decision string `json:"decision"`
Reason string `json:"reason,omitempty"`
}

func (h *Handler) PostPermissionDecide(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
errcode.Write(w, errcode.ErrMethodNotAllowed)
return
}

var req struct {
RunID string `json:"runId"`
RequestID string `json:"requestId"`
Decision string `json:"decision"`
Reason string `json:"reason,omitempty"`
}
var req permissionDecideRequest
if err := decodeOptionalJSON(r, &req); err != nil {
errcode.Write(w, errcode.ErrInvalidJSON)
return
}
req.RunID = strings.TrimSpace(req.RunID)
req.RequestID = strings.TrimSpace(req.RequestID)
req.Decision = strings.TrimSpace(req.Decision)
req.ControlID = strings.TrimSpace(req.ControlID)
req.HubTaskID = strings.TrimSpace(req.HubTaskID)
if req.RunID == "" {
errcode.Write(w, errcode.ErrRunIDRequired)
return
Expand All @@ -42,51 +51,117 @@ func (h *Handler) PostPermissionDecide(w http.ResponseWriter, r *http.Request) {
errcode.Write(w, errcode.ErrInvalidDecision)
return
}
if (req.ControlID == "") != (req.HubTaskID == "") {
errcode.Write(w, errcode.ErrBadRequest.WithMessage("controlId and hubTaskId must be provided together"))
return
}

// Ownership gate, resolved before any state change (broker decide / registry
// consume / event publish): without it any caller past the Edge's coarse auth
// could allow or deny a live tool call of somebody else's run, i.e. make the
// victim's agent perform the call. The 404 body is byte-identical to the
// "no such request" path below (same errcode, no distinguishing message) so a
// foreign runId and a nonexistent runId stay indistinguishable — this endpoint
// must not become a runId existence oracle. Local single-tenant mode resolves
// to the documented bypass sentinel and is unaffected; an empty principal under
// Hub JWT fails closed (AH-SR-045).
if !isRunOwnedBy(ensureStore(h), req.RunID, h.ownerUserID(r)) {
// Ownership and task binding are resolved before any state change:
// broker decide, registry consume, event publish, and receipt storage all
// happen only for the real run owner and the exact stored Hub task. The
// 404 body is byte-identical for a foreign runId, a missing runId, a
// wrong hubTaskId, and a missing pending request so this endpoint is not an
// existence oracle. Local single-tenant mode resolves to the documented
// bypass sentinel; empty principal under Hub JWT fails closed (AH-SR-045).
repo := ensureStore(h)
owner := h.ownerUserID(r)
if !isRunOwnedBy(repo, req.RunID, owner) {
errcode.Write(w, errcode.ErrPermissionRequestNotFound)
return
}
if req.ControlID != "" {
run, ok := repo.GetRun(req.RunID)
if !ok || run.HubTaskID != req.HubTaskID {
errcode.Write(w, errcode.ErrPermissionRequestNotFound)
return
}
}

registry := h.ensurePermissionRegistry()
permission, ok := pendingPermissionFromBroker(h.ensurePermissionBroker(), req.RunID, req.RequestID, req.Decision, req.Reason)
if ok {
_, _ = registry.Consume(req.RunID, req.RequestID)
} else {
permission, ok = registry.Consume(req.RunID, req.RequestID)
if req.ControlID == "" {
h.decideLegacyPermission(w, registry, req)
return
}
h.decideModernPermission(w, registry, req)
}

func (h *Handler) decideLegacyPermission(w http.ResponseWriter, registry *permission.PermissionRegistry, req permissionDecideRequest) {
pending, ok := h.resolvePendingPermission(registry, req)
if !ok {
errcode.Write(w, errcode.ErrPermissionRequestNotFound)
return
}
h.publishPermissionDecision(pending, req)
slog.Info("permission decided by Desktop", "requestId", req.RequestID, "decision", req.Decision)
writeSuccess(w, http.StatusOK, map[string]any{"status": "ok"})
}

func (h *Handler) decideModernPermission(w http.ResponseWriter, registry *permission.PermissionRegistry, req permissionDecideRequest) {
receipt := permissionReceipt{
ControlID: req.ControlID,
HubTaskID: req.HubTaskID,
RunID: req.RunID,
RequestID: req.RequestID,
Decision: req.Decision,
Reason: req.Reason,
}
result := h.ensurePermissionReceipts().apply(receipt, func() bool {
pending, ok := h.resolvePendingPermission(registry, req)
if !ok {
errcode.Write(w, errcode.ErrPermissionRequestNotFound)
return
return false
}
h.publishPermissionDecision(pending, req)
return true
})
switch {
case result.Conflict:
errcode.Write(w, errcode.ErrConflict.WithMessage("permission decision conflict"))
return
case result.Full:
errcode.Write(w, errcode.ErrTooManyRequests.WithMessage("permission decision receipt cache is full"))
return
case !result.Applied:
errcode.Write(w, errcode.ErrPermissionRequestNotFound)
return
}
slog.Info("permission decision applied by Desktop", "requestId", req.RequestID, "controlId", req.ControlID, "decision", req.Decision)
writeSuccess(w, http.StatusOK, map[string]any{
"status": "ok",
"controlId": result.Receipt.ControlID,
"hubTaskId": result.Receipt.HubTaskID,
"runId": result.Receipt.RunID,
"requestId": result.Receipt.RequestID,
"decision": result.Receipt.Decision,
"applied": true,
"deduplicated": result.Deduplicated,
})
}

func (h *Handler) resolvePendingPermission(registry *permission.PermissionRegistry, req permissionDecideRequest) (permission.PendingPermission, bool) {
pending, ok := pendingPermissionFromBroker(h.ensurePermissionBroker(), req.RunID, req.RequestID, req.Decision, req.Reason)
if ok {
_, _ = registry.Consume(req.RunID, req.RequestID)
return pending, true
}
return registry.Consume(req.RunID, req.RequestID)
}

scope := map[string]any{"runId": permission.RunID}
if permission.ProjectID != "" {
scope["projectId"] = permission.ProjectID
func (h *Handler) publishPermissionDecision(pending permission.PendingPermission, req permissionDecideRequest) {
scope := map[string]any{"runId": pending.RunID}
if pending.ProjectID != "" {
scope["projectId"] = pending.ProjectID
}
if permission.ThreadID != "" {
scope["threadId"] = permission.ThreadID
if pending.ThreadID != "" {
scope["threadId"] = pending.ThreadID
}
ensureBus(h).Publish(adapters.BusEventPermissionDecided, scope, map[string]any{
"runId": req.RunID,
"requestId": req.RequestID,
"toolName": permission.ToolName,
"toolUseId": permission.ToolUseID,
"toolName": pending.ToolName,
"toolUseId": pending.ToolUseID,
"decision": req.Decision,
"reason": req.Reason,
})

slog.Info("permission decided by Desktop", "requestId", req.RequestID, "decision", req.Decision)
writeSuccess(w, http.StatusOK, map[string]any{"status": "ok"})
}

func pendingPermissionFromBroker(broker *adapters.PermissionDecisionBroker, runID, requestID, decision, reason string) (permission.PendingPermission, bool) {
Expand Down
5 changes: 3 additions & 2 deletions edge-server/internal/api/handlers_run_callback.go
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,8 @@ func validateReplayCallbackOwner(req runRequest, run store.Run) *errcode.Error {

func (h *Handler) runCallbackCapabilities() map[string]bool {
return map[string]bool{
"runCallbackOwnership": true,
"directHubCallbacks": h.directHubCallbacksConfigured(),
"runCallbackOwnership": true,
"directHubCallbacks": h.directHubCallbacksConfigured(),
"permissionDecisionReceipts": true,
}
}
183 changes: 183 additions & 0 deletions edge-server/internal/api/permission_receipts.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,183 @@
package api

import (
"strings"
"sync"
"time"

"github.com/agenthub/edge-server/internal/deliverydedup"
)

const (
permissionReceiptDefaultCapacity = deliverydedup.DefaultCapacity
permissionReceiptDefaultTTL = deliverydedup.DefaultTTL
)

// permissionRunRequestKey is the business identity that may accept only one
// modern control. A different control ID for the same run/request is a
// conflict, never a second application.
type permissionRunRequestKey struct {
runID string
requestID string
}

// permissionReceipt is an applied modern permission decision. It is a warm,
// bounded replay receipt only; it is never authority and never a pending
// request.
type permissionReceipt struct {
ControlID string
HubTaskID string
RunID string
RequestID string
Decision string
Reason string
expiresAt time.Time
}

func (r permissionReceipt) sameAs(other permissionReceipt) bool {
return r.ControlID == other.ControlID &&
r.HubTaskID == other.HubTaskID &&
r.RunID == other.RunID &&
r.RequestID == other.RequestID &&
r.Decision == other.Decision &&
r.Reason == other.Reason
}

func (r permissionReceipt) runRequestKey() permissionRunRequestKey {
return permissionRunRequestKey{runID: r.RunID, requestID: r.RequestID}
}

// permissionReceiptApplyResult reports the outcome of an atomic modern
// decision application. Full and Conflict are returned before the callback
// runs, so no broker/consume/event effect happens after those failures.
type permissionReceiptApplyResult struct {
Receipt permissionReceipt
Applied bool
Deduplicated bool
Conflict bool
Full bool
}

// permissionReceiptCache is a process-local, bounded TTL cache of applied
// modern permission receipts. It deliberately has no durable state and does
// not retain pending requests: a cold miss must fall through to the existing
// broker/registry and return not-found instead of inventing success. Capacity
// pressure rejects a new decision before effects; it does not evict an applied
// receipt to admit another control.
type permissionReceiptCache struct {
mu sync.Mutex
capacity int
ttl time.Duration
now func() time.Time
byControl map[string]permissionReceipt
byRunReq map[permissionRunRequestKey]struct{}
}

func newPermissionReceiptCache(capacity int, ttl time.Duration) *permissionReceiptCache {
if capacity <= 0 {
panic("api: permission receipt cache capacity must be > 0")
}
if ttl <= 0 {
panic("api: permission receipt cache ttl must be > 0")
}
return &permissionReceiptCache{
capacity: capacity,
ttl: ttl,
now: time.Now,
byControl: make(map[string]permissionReceipt, capacity),
byRunReq: make(map[permissionRunRequestKey]struct{}, capacity),
}
}

func (c *permissionReceiptCache) withClock(clock func() time.Time) *permissionReceiptCache {
c.mu.Lock()
defer c.mu.Unlock()
c.now = clock
return c
}

// apply atomically decides whether a modern request is a warm replay, a
// conflict, a capacity rejection, or a newly pending application. The callback
// runs while the cache lock is held so one winner performs every effect and
// concurrent retries observe a stable receipt afterwards.
func (c *permissionReceiptCache) apply(request permissionReceipt, applyFn func() bool) permissionReceiptApplyResult {
c.mu.Lock()
defer c.mu.Unlock()
request.ControlID = strings.TrimSpace(request.ControlID)
request.HubTaskID = strings.TrimSpace(request.HubTaskID)
request.RunID = strings.TrimSpace(request.RunID)
request.RequestID = strings.TrimSpace(request.RequestID)
request.Decision = strings.TrimSpace(request.Decision)
if request.ControlID == "" || request.HubTaskID == "" || request.RunID == "" || request.RequestID == "" {
return permissionReceiptApplyResult{}
}
c.purgeExpiredLocked(c.now())

if stored, ok := c.byControl[request.ControlID]; ok {
if stored.sameAs(request) {
return permissionReceiptApplyResult{
Receipt: stored,
Applied: true,
Deduplicated: true,
}
}
return permissionReceiptApplyResult{Conflict: true}
}
if _, ok := c.byRunReq[request.runRequestKey()]; ok {
return permissionReceiptApplyResult{Conflict: true}
}
if len(c.byControl) >= c.capacity {
return permissionReceiptApplyResult{Full: true}
}
if !applyFn() {
return permissionReceiptApplyResult{}
}
request.expiresAt = c.now().Add(c.ttl)
c.storeLocked(request)
return permissionReceiptApplyResult{
Receipt: request,
Applied: true,
}
}

func (c *permissionReceiptCache) storeLocked(receipt permissionReceipt) {
receipt.ControlID = strings.TrimSpace(receipt.ControlID)
receipt.HubTaskID = strings.TrimSpace(receipt.HubTaskID)
receipt.RunID = strings.TrimSpace(receipt.RunID)
receipt.RequestID = strings.TrimSpace(receipt.RequestID)
receipt.Decision = strings.TrimSpace(receipt.Decision)
controlID := receipt.ControlID
if controlID == "" {
return
}
c.byControl[controlID] = receipt
c.byRunReq[receipt.runRequestKey()] = struct{}{}
}

func (c *permissionReceiptCache) purgeExpiredLocked(now time.Time) {
for id, receipt := range c.byControl {
if now.After(receipt.expiresAt) {
delete(c.byControl, id)
delete(c.byRunReq, receipt.runRequestKey())
}
}
}

func (c *permissionReceiptCache) len() int {
c.mu.Lock()
defer c.mu.Unlock()
c.purgeExpiredLocked(c.now())
return len(c.byControl)
}

func (h *Handler) ensurePermissionReceipts() *permissionReceiptCache {
h.permissionRegistryMu.Lock()
defer h.permissionRegistryMu.Unlock()
if h.permissionReceipts == nil {
h.permissionReceipts = newPermissionReceiptCache(
permissionReceiptDefaultCapacity,
permissionReceiptDefaultTTL,
)
}
return h.permissionReceipts
}
Loading
Loading