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
6 changes: 3 additions & 3 deletions services/core/IMPLEMENTATION.md

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -25,10 +25,7 @@ func TestNativeClassificationPostgresRoundTripAndPublicPrivacy(t *testing.T) {
if err != nil {
t.Fatal(err)
}
receipt, err := s.SubmitMessage(t.Context(), tenant, session.ID, "input", json.RawMessage(`{"text":"test"}`))
if err != nil {
t.Fatal(err)
}
receipt := submitMessage(t, pool, tenant, session.ID, "input", json.RawMessage(`{"text":"test"}`))
transitionTurn(t, pool, tenant, session.ID, receipt.TurnID, sessions.TurnTransition{ExpectedStatus: sessions.TurnQueued, Status: sessions.TurnInProgress})
status := 503
result := execution.Result{ErrorCode: "engine_failed", Error: "Bearer secret-canary https://private.example/key", EngineErrorCode: code, EngineHTTPStatus: &status, Done: proto.DonePayload{Usage: proto.Usage{InputTokens: 7, OutputTokens: 3}, Metadata: map[string]any{proto.DoneMetaAgentSessionID: "native-secret-canary"}}}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,20 @@ func databaseSessionReads(pool *pgxpool.Pool) func(*Dependencies, *testFakes) {
}
}

// submitMessage admits one message input through the Session service on pool.
func submitMessage(t *testing.T, pool *pgxpool.Pool, tenant, session, key string, payload json.RawMessage) sessions.InputReceipt {
t.Helper()
service, err := sessions.NewService(sessionpg.New(pgunit.NewPool(pool), nil))
if err != nil {
t.Fatal(err)
}
receipts, err := service.SubmitInputs(t.Context(), tenant, session, key, []sessions.Input{{Kind: "message", Payload: payload}})
if err != nil {
t.Fatal(err)
}
return receipts[0]
}

// transitionTurn moves the Turn as the execution owner does, over a pooled
// Session transaction.
func transitionTurn(t *testing.T, pool *pgxpool.Pool, tenant, session, turn string, transition sessions.TurnTransition) {
Expand Down
5 changes: 1 addition & 4 deletions services/core/internal/api/session_diagnostics_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,10 +19,7 @@ func TestDiagnosticsCoreHandlerDatabaseBoundary(t *testing.T) {
if err != nil {
t.Fatal(err)
}
receipt, err := s.SubmitMessage(t.Context(), tenant, session.ID, "input", json.RawMessage(`{"text":"input-secret-canary"}`))
if err != nil {
t.Fatal(err)
}
receipt := submitMessage(t, pool, tenant, session.ID, "input", json.RawMessage(`{"text":"input-secret-canary"}`))
transitionTurn(t, pool, tenant, session.ID, receipt.TurnID, sessions.TurnTransition{ExpectedStatus: sessions.TurnQueued, Status: sessions.TurnInProgress})
transitionTurn(t, pool, tenant, session.ID, receipt.TurnID, sessions.TurnTransition{ExpectedStatus: sessions.TurnInProgress, Status: sessions.TurnFailed, Outcome: json.RawMessage(`{"error_code":"device_disconnected","error":"Bearer raw-secret-canary https://private.example/key","done":{"native_id":"secret-native-canary"}}`)})
base := adminSessionsPath + session.ID
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -92,10 +92,12 @@ func TestArchiveWaitingCleanupReceiptBarrier(t *testing.T) {
if err != nil {
t.Fatal(err)
}
input, err := s.SubmitMessage(t.Context(), project.TenantID, session.ID, "start", json.RawMessage(`{"text":"run"}`))
_, service := testSessions(t, pool, nil)
inputs, err := service.SubmitInputs(t.Context(), project.TenantID, session.ID, "start", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"run"}`)}})
if err != nil {
t.Fatal(err)
}
input := inputs[0]
// This fixture isolates lifecycle ordering. Protocol-driven waiting is
// independently exercised in TestArchiveWaitingCancellationReceipts.
for _, transition := range []sessions.TurnTransition{{ExpectedStatus: sessions.TurnQueued, Status: sessions.TurnInProgress}, {ExpectedStatus: sessions.TurnInProgress, Status: sessions.TurnWaiting}} {
Expand Down
2 changes: 1 addition & 1 deletion services/core/internal/execution/delivery.go
Original file line number Diff line number Diff line change
Expand Up @@ -303,7 +303,7 @@ func (d *Dispatcher) deliver(ctx context.Context, tenantID, sessionID string, pe
continue
}
if pending == nil {
inputs, err := d.Store.ListTurnInputs(ctx, tenantID, sessionID, request.RunID, result.AppliedThrough, 1)
inputs, err := d.SessionsReader.ListTurnInputs(ctx, tenantID, sessionID, request.RunID, result.AppliedThrough, 1)
if err != nil {
result.ErrorCode = "execution_state_unavailable"
if ctx.Err() != nil {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -82,10 +82,12 @@ func newFinishObservationFixture(t *testing.T, maxConnections int32) finishObser
}
func (f finishObservationFixture) start(t *testing.T) sessions.InputReceipt {
t.Helper()
receipt, err := f.s.SubmitMessage(t.Context(), f.tenant, f.session.ID, uuid.NewString(), json.RawMessage(`{"input":[{"role":"user","content":[{"type":"input_text","text":"fixture"}]}]}`))
_, service := testSessions(t, f.pool, nil)
receipts, err := service.SubmitInputs(t.Context(), f.tenant, f.session.ID, uuid.NewString(), []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"input":[{"role":"user","content":[{"type":"input_text","text":"fixture"}]}]}`)}})
if err != nil {
t.Fatal(err)
}
receipt := receipts[0]
if _, err = f.execution.TransitionTurn(t.Context(), f.tenant, f.session.ID, receipt.TurnID, sessions.TurnTransition{ExpectedStatus: sessions.TurnQueued, Status: sessions.TurnInProgress}); err != nil {
t.Fatal(err)
}
Expand Down
6 changes: 3 additions & 3 deletions services/core/internal/execution/environment_admission.go
Original file line number Diff line number Diff line change
Expand Up @@ -97,7 +97,7 @@ func (w *Worker) submitEnvironmentInputs(ctx context.Context, session sessions.S
changed, unsubscribe := w.dispatcher.notifications.subscribe(session.TenantID, session.ID)
defer unsubscribe()
reserve, cancel := context.WithTimeout(ctx, 5*time.Second)
reservation, err := w.admission.ReserveEnvironmentInput(reserve, session.TenantID, session.ID, key, inputs)
reservation, err := w.dispatcher.Sessions.ReserveEnvironmentInput(reserve, session.TenantID, session.ID, key, inputs)
cancel()
if err != nil {
return nil, err
Expand Down Expand Up @@ -140,9 +140,9 @@ func (w *Worker) environmentInputOutcome(ctx context.Context, session sessions.S
defer cancel()
// The database rechecks its clock under the Session lock before settlement.
if !time.Now().Before(reservation.Deadline) {
return w.admission.ExpireEnvironmentInput(read, session.TenantID, session.ID, reservation.ID)
return w.dispatcher.Sessions.ExpireEnvironmentInput(read, session.TenantID, session.ID, reservation.ID)
}
return w.admission.GetEnvironmentInputReservation(read, session.TenantID, session.ID, reservation.ID)
return w.dispatcher.SessionsReader.GetEnvironmentInputReservation(read, session.TenantID, session.ID, reservation.ID)
}

func (w *Worker) checkAdmissionOwnership(ctx context.Context) error {
Expand Down
2 changes: 1 addition & 1 deletion services/core/internal/execution/message_input.go
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ func messageInput(raw json.RawMessage) (proto.MessageInput, error) {
}

func (d *Dispatcher) initialInput(ctx context.Context, tenant, session, turn string) (proto.MessageInput, int64, error) {
inputs, err := d.Store.ListTurnInputs(ctx, tenant, session, turn, 0, 100)
inputs, err := d.SessionsReader.ListTurnInputs(ctx, tenant, session, turn, 0, 100)
if err != nil {
return nil, 0, err
}
Expand Down
4 changes: 3 additions & 1 deletion services/core/internal/execution/preparation.go
Original file line number Diff line number Diff line change
Expand Up @@ -97,7 +97,9 @@ func (d *Dispatcher) awaitPreparation(ctx context.Context, tenant, session strin
case <-ctx.Done():
return pending, ctx.Err()
case <-tick.C:
current, err := d.Store.ExpireEnvironmentInput(ctx, tenant, session, pending.ID)
expire, cancel := context.WithTimeout(ctx, 5*time.Second)
current, err := d.Sessions.ExpireEnvironmentInput(expire, tenant, session, pending.ID)
cancel()
if err != nil || current.State != sessions.EnvironmentInputPending {
return current, err
}
Expand Down
10 changes: 7 additions & 3 deletions services/core/internal/execution/prepared_dispatch.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import (
"encoding/json"
"errors"
"strings"
"time"

v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1"
"github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto"
Expand All @@ -17,12 +18,15 @@ type EnvironmentRun struct {
}

// RunEnvironmentInput reserves a Turn on the Session-owned Runtime Executor. It
// checks lease, the lease d.Store was built on, before any Runtime preparation.
// checks lease, the lease the Dispatcher's execution operations hold, before
// any Runtime preparation.
func (d *Dispatcher) RunEnvironmentInput(ctx context.Context, lease Ownership, tenantID, sessionID, reservationID string) (run EnvironmentRun, err error) {
if err = lease.CheckOwnership(ctx); err != nil {
return run, err
}
run.Reservation, err = d.Store.ExpireEnvironmentInput(ctx, tenantID, sessionID, reservationID)
expire, cancel := context.WithTimeout(ctx, 5*time.Second)
run.Reservation, err = d.Sessions.ExpireEnvironmentInput(expire, tenantID, sessionID, reservationID)
cancel()
if err != nil || run.Reservation.State != sessions.EnvironmentInputPending {
return run, err
}
Expand Down Expand Up @@ -92,7 +96,7 @@ func (d *Dispatcher) RunEnvironmentInput(ctx context.Context, lease Ownership, t
if err := d.messageInputSupport(peer, session.Engine, snapshot, messages); err != nil {
return run, err
}
promoted, err := d.Store.PromoteEnvironmentInput(owner, tenantID, sessionID, reservationID)
promoted, err := d.sessionExecution.PromoteEnvironmentInput(owner, tenantID, sessionID, reservationID)
if errors.Is(err, sessions.ErrTurnConflict) {
// A rejected claim leaves the reservation pending for a later attempt.
return run, err
Expand Down
2 changes: 1 addition & 1 deletion services/core/internal/execution/worker.go
Original file line number Diff line number Diff line change
Expand Up @@ -303,7 +303,7 @@ func (w *Worker) Run(ctx context.Context) (runErr error) {
return err
}
if maintenance {
if _, err := w.dispatcher.Store.ExpireEnvironmentInputs(ctx); err != nil {
if _, err := w.dispatcher.sessionExecution.ExpireEnvironmentInputs(ctx); err != nil {
w.observeSchedulerPoll(0, err)
return err
}
Expand Down
8 changes: 4 additions & 4 deletions services/core/internal/execution/worker_schedule.go
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ func (s *workerSchedule) selectWork(ctx context.Context, w *Worker, devices []st
}
var environments []sessions.EnvironmentInputWork
if !time.Now().Before(s.nextEnvironmentScan) {
environments, err = w.dispatcher.Store.ListEnvironmentInputWork(ctx, s.environmentCursor, devices)
environments, err = w.dispatcher.SessionsReader.ListEnvironmentInputWork(ctx, s.environmentCursor, devices)
if err != nil {
return nil, err
}
Expand All @@ -43,7 +43,7 @@ func (s *workerSchedule) selectWork(ctx context.Context, w *Worker, devices []st
s.environmentCursor = ""
// Retry the first page now instead of spending a scan interval on EOF.
// A single refill preserves the candidate bound and cannot spin when empty.
environments, err = w.dispatcher.Store.ListEnvironmentInputWork(ctx, "", devices)
environments, err = w.dispatcher.SessionsReader.ListEnvironmentInputWork(ctx, "", devices)
if err != nil {
return nil, err
}
Expand Down Expand Up @@ -94,10 +94,10 @@ func (w *Worker) runEnvironmentInput(ctx context.Context, item scheduledWork) er
return nil
}
if errors.Is(err, ErrModelProviderRequired) {
return w.dispatcher.Store.FailEnvironmentInput(ctx, item.TenantID, item.SessionID, item.reservationID, "model_provider_required")
return w.dispatcher.sessionExecution.FailEnvironmentInput(ctx, item.TenantID, item.SessionID, item.reservationID, "model_provider_required")
}
if errors.Is(err, errPreparationFailed) && run.Reservation.State == sessions.EnvironmentInputPending {
return w.dispatcher.Store.FailEnvironmentInput(ctx, item.TenantID, item.SessionID, item.reservationID, "runtime_preparation_failed")
return w.dispatcher.sessionExecution.FailEnvironmentInput(ctx, item.TenantID, item.SessionID, item.reservationID, "runtime_preparation_failed")
}
if run.Reservation.State == sessions.EnvironmentInputAdmitted {
return err
Expand Down
2 changes: 1 addition & 1 deletion services/core/internal/execution/worker_wakeup.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ func (w *Worker) wakeScheduler() {
}

func (w *Worker) admitInputs(ctx context.Context, tenant, session, key string, inputs []sessions.Input) ([]sessions.InputReceipt, error) {
receipts, err := w.admission.SubmitInputs(ctx, tenant, session, key, inputs)
receipts, err := w.dispatcher.Sessions.SubmitInputs(ctx, tenant, session, key, inputs)
if err == nil {
w.wakeScheduler()
w.dispatcher.notifications.notify(tenant, session)
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,165 @@
package sessionpg

import (
"context"
"encoding/json"
"errors"
"reflect"
"sync"
"testing"
"time"

"github.com/jackc/pgx/v5/pgtype"
"github.com/jackc/pgx/v5/pgxpool"

"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/db/sqlc"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgtest"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgunit"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions"
)

var reservedBatch = []sessions.Input{messageInput("first"), messageInput("second")}

// environmentInputHistory checks the Session's Turns and admitted inputs, one
// Item per input, and that no input left Turn or Item events.
func environmentInputHistory(t *testing.T, pool *pgxpool.Pool, session pgtype.UUID, turns, inputs int) {
t.Helper()
var gotTurns, gotInputs, items, events int
err := pool.QueryRow(t.Context(), `
SELECT (SELECT count(*) FROM turns WHERE session_id=$1),
(SELECT count(*) FROM turn_inputs WHERE session_id=$1),
(SELECT count(*) FROM session_items WHERE session_id=$1),
(SELECT count(*) FROM session_events WHERE session_id=$1
AND (payload ? 'turn' OR payload->'event' ? 'item'))`, session).Scan(&gotTurns, &gotInputs, &items, &events)
if err != nil || gotTurns != turns || gotInputs != inputs || items != inputs || (inputs == 0 && events != 0) {
t.Fatal("history", gotTurns, gotInputs, items, events, err)
}
}

// reserve reserves the two-message batch under key.
func reserve(t *testing.T, service *sessions.Service, tenant, session pgtype.UUID, key string) sessions.EnvironmentInputReservation {
t.Helper()
got, err := service.ReserveEnvironmentInput(t.Context(), text(tenant), text(session), key, reservedBatch)
if err != nil {
t.Fatal(err)
}
return got
}

// passDeadline moves the reservation's deadline into the past.
func passDeadline(t *testing.T, pool *pgxpool.Pool, reservation string) {
t.Helper()
exec(t, pool, `UPDATE environment_input_reservations SET deadline=clock_timestamp()-interval '1 second' WHERE id=$1`, reservation)
}

// cancelPending cancels the Session's pending input as Session cancellation
// does, then reads the reservation back.
func cancelPending(ctx context.Context, pool *pgxpool.Pool, store *Store, tenant, session pgtype.UUID, reservation string) (sessions.EnvironmentInputReservation, error) {
err := WithSession(ctx, pgunit.NewPool(pool), tenant, session, func(ctx context.Context, q *sqlc.Queries, locked sessions.LockedSession) error {
if err := locked.Public(); err != nil {
return err
}
bound := BindSession(q, tenant, session)
return sessions.TrackInputActivity(ctx, bound, bound.CancelPendingInput)
})
if err != nil {
return sessions.EnvironmentInputReservation{}, err
}
return store.GetEnvironmentInputReservation(ctx, text(tenant), text(session), reservation)
}

// Concurrent equivalent reservations from two pools share one pending
// identity and deadline, which a restart keeps.
func TestEnvironmentInputReservationConcurrentIdentity(t *testing.T) {
pool := pgtest.Open(t)
_, service := stagingService(t, pool)
_, other := stagingService(t, pgtest.Open(t))
tenant, session, _ := newEnvironment(t, pool, "self_hosted", "pending")
ctx := t.Context()
batch := []sessions.Input{
{Kind: "message", Payload: json.RawMessage(`{"text":"first","detail":{"a":1,"b":2}}`)},
messageInput("second"),
}
const count = 8
results := make(chan sessions.EnvironmentInputReservation, count)
var wg sync.WaitGroup
for i := range count {
wg.Go(func() {
admission := service
inputs := append([]sessions.Input(nil), batch...)
if i%2 == 0 {
admission = other
inputs[0].Payload = json.RawMessage(` { "detail": {"b": 2, "a": 1}, "text": "first" } `)
}
got, err := admission.ReserveEnvironmentInput(ctx, text(tenant), text(session), "request", inputs)
if err != nil {
t.Error(err)
return
}
results <- got
})
}
wg.Wait()
close(results)
var first sessions.EnvironmentInputReservation
received := 0
for result := range results {
received++
if first.ID == "" {
first = result
}
if !reflect.DeepEqual(first, result) {
t.Fatal("reservation identity changed", first, result)
}
}
if received != count || first.State != sessions.EnvironmentInputPending || first.ID == "" || first.Deadline.Sub(first.CreatedAt) != 5*time.Minute || first.SettledAt != nil || len(first.Receipts) != 0 {
t.Fatal("invalid pending result", received, first)
}
environmentInputHistory(t, pool, session, 0, 0)
for _, changed := range [][]sessions.Input{batch[:1], {batch[1], batch[0]}, {messageInput("changed"), batch[1]}} {
if _, err := service.ReserveEnvironmentInput(ctx, text(tenant), text(session), "request", changed); !errors.Is(err, sessions.ErrIdempotencyConflict) {
t.Fatal("changed request accepted", err)
}
}
if _, err := other.ReserveEnvironmentInput(ctx, text(tenant), text(session), "other", batch); !errors.Is(err, sessions.ErrTurnConflict) {
t.Fatal("second pending request accepted", err)
}
pool.Close()
_, restarted := stagingService(t, pgtest.Open(t))
got, err := restarted.ReserveEnvironmentInput(ctx, text(tenant), text(session), "request", batch)
if err != nil || !reflect.DeepEqual(first, got) {
t.Fatal("restart changed deadline or identity", got, err)
}
}

// A key that direct admission already used keeps its receipts and gains no
// reservation, and its retry leaves newer pending input alone.
func TestEnvironmentInputReservationKeepsEarlierDirectIdentity(t *testing.T) {
pool := pgtest.Open(t)
store, service := stagingService(t, pool)
tenant, session, _ := newEnvironment(t, pool, "self_hosted", "pending")
ctx := t.Context()
input := messageInput("already admitted")
receipts, err := service.SubmitInputs(ctx, text(tenant), text(session), "direct", []sessions.Input{input})
if err != nil {
t.Fatal(err)
}
move(t, pool, tenant, session, receipts[0].TurnID, sessions.TurnQueued, sessions.TurnInProgress)
move(t, pool, tenant, session, receipts[0].TurnID, sessions.TurnInProgress, sessions.TurnCompleted)
pending := reserve(t, service, tenant, session, "new")
got, err := service.ReserveEnvironmentInput(ctx, text(tenant), text(session), "direct", []sessions.Input{input})
if err != nil || got.State != sessions.EnvironmentInputAdmitted || got.ID != "" || !got.Deadline.IsZero() || len(got.Receipts) != 1 || got.Receipts[0].Sequence != receipts[0].Sequence {
t.Fatal("direct admission gained a reservation", got, err)
}
if _, err := service.ReserveEnvironmentInput(ctx, text(tenant), text(session), "direct", []sessions.Input{messageInput("changed")}); !errors.Is(err, sessions.ErrIdempotencyConflict) {
t.Fatal(err)
}
retry, err := service.SubmitInputs(ctx, text(tenant), text(session), "direct", []sessions.Input{input})
if err != nil || len(retry) != 1 || !retry[0].Replayed {
t.Fatal(retry, err)
}
retained, err := store.GetEnvironmentInputReservation(ctx, text(tenant), text(session), pending.ID)
if err != nil || !reflect.DeepEqual(retained, pending) {
t.Fatal("old retry affected new pending input", retained, err)
}
}
Loading
Loading