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
12 changes: 7 additions & 5 deletions services/core/IMPLEMENTATION.md
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,8 @@ Two errors are shared across domains, each with one `api` helper: `textvalue.Err

Shared vocabulary has one owner each, and domains use it rather than copy it. `internal/environmentconfig` owns Environment setup, Skills, Plugins and initial files with their validation and public metadata; `Setup.Validate` checks requested configuration, where a Skill may be an unresolved reference, and `Setup.ValidateInstalled` checks frozen, installable configuration. `internal/skills` owns `ParseVersion`, the canonical positive decimal Skill version. `internal/metadata` owns the metadata rules: `Validate` for the pair, key and value limits and U+0000, `ValidateStorable` for U+0000 alone, and `Encode` with its 64 KiB bound. `internal/jsonobject` owns `Normalize`, the stable encoding of stored JSON objects that snapshots and retry identities compare. These packages import no persistence.

`internal/sessions` owns the Session change vocabulary and its decisions: the public changes that report Turn and Session transitions, what a Turn that ends settles, measured Turn usage and the Session activity a change reports. `internal/items` owns public Items: it projects observations, merges them into stored Items and builds the ordered events that report each Item change. Neither imports persistence. `internal/persistence/postgres/sessionpg` loads the facts those decisions read and applies them inside the caller's Session transaction, under the Session lock: it allocates event sequence positions, event IDs, Item positions and output indexes, writes the journal, Items, Turn usage and Artifact settlement, and prunes the journal. It decides nothing.

`store` is transitional. `store.New` builds a pooled Store, and `store.NewExecution(s, lease)` builds the execution writer on a lease it borrows. An execution-only operation on a pooled Store fails with `store.ErrExecutionAuthority`. New adapters do not copy that check: their execution repositories require a `*pgunit.Lease` at construction, their public repositories expose no execution operation, and the check goes away with `store`.

`cmd/server` owns the execution lease. It acquires one `pgunit.Lease`, builds every lease-bound adapter on it, and passes the lease and those adapters together as one `execution.Owner` to `execution.StartWorker`. If anything fails before that call, `cmd/server` closes the lease. From that call the Worker owns cleanup: a failed start closes the lease before it returns, and a started Worker closes it after `Run` has cancelled and drained its work. Each close runs under its own bounded deadline, independent of the cancelled request or run. Lease-bound adapters and the store writer borrow the lease and never close it, and the Worker uses the lease only through `Owner.Lease`, never through an adapter. Store integration tests start the Worker the same way through `startWorker`.
Expand Down Expand Up @@ -44,7 +46,7 @@ Source Files are Project resources with a lifecycle independent of copied worksp

Session Artifacts are immutable published copies, separate from live workspace files and source Files. The private output exporter reuses the authorized workspace path boundary and streams bounded bytes; publication requires complete capture and confirmed helper and transport success, not merely valid archive syntax. The daemon owns and drains the exporter's stdout pipe separately from child reaping, so pull-transport backpressure cannot consume the process-exit I/O deadline; after helper exit, each pipe read has one second, reset after consumer delays, which rejects inherited pipes that never close. Cancellation closes the owned reader and the dispatch consumer, then waits for the child. Never extract an output archive into Core's filesystem or hold the execution lease through a large transfer.

Capture bytes into private large objects without a Session admission lock. Before capture, seal native input under that lock with the private capture marker; later messages reuse the Environment input reservation and wait for the next Turn, and the public Turn stays in progress until publication settles. Directory reads during capture use an independent read-only preparation, not the released native Run. After a confirmed export, lock and recheck the live Turn, then publish the metadata in the same transaction as the Turn's completion. Decide Artifact republication in the Turn's terminal transaction, not the capture transaction: capture commits and releases the Session lock before the Turn completes, and the terminal transaction holds the Session lock that also orders Artifact deletion. Drop staged rows whose SHA-256 equals the newest remaining published Artifact for the path, ordered by the producing Turn's creation time and then ID (publication time can come from the Runtime and does not order Turns), and unlink their large objects in the same transaction. Published rows are never modified. Failed and cancelled Turns discard private objects, and Session deletion removes private and published copies. The exporter skips output symlinks by their `lstat` type without following them; hard links, other special files, device crossings and concurrent changes reject the capture. Stored reads are authorized independently of Environment availability, so published Artifacts outlive the Environment.
Capture bytes into private large objects without a Session admission lock. Before capture, seal native input under that lock with the private capture marker; later messages reuse the Environment input reservation and wait for the next Turn, and the public Turn stays in progress until publication settles. Directory reads during capture use an independent read-only preparation, not the released native Run. After a confirmed export, lock and recheck the live Turn, then publish the metadata in the same transaction as the Turn's completion. `sessions.EndTurn` decides Artifact publication and `sessionpg` applies it in the Turn's terminal transaction, not the capture transaction: capture commits and releases the Session lock before the Turn completes, and the terminal transaction holds the Session lock that also orders Artifact deletion. Drop staged rows whose SHA-256 equals the newest remaining published Artifact for the path, ordered by the producing Turn's creation time and then ID (publication time can come from the Runtime and does not order Turns), and unlink their large objects in the same transaction. Published rows are never modified. Failed and cancelled Turns discard private objects, and Session deletion removes private and published copies. The exporter skips output symlinks by their `lstat` type without following them; hard links, other special files, device crossings and concurrent changes reject the capture. Stored reads are authorized independently of Environment availability, so published Artifacts outlive the Environment.

- A malformed `environment_id` Artifact filter resolves to the never-assigned maximum UUID, so it matches nothing without a text comparison.

Expand Down Expand Up @@ -139,7 +141,7 @@ Accepted public MCP credential profiles are declared centrally by the service, s

`services/core/internal/execution` claims a Turn from `queued` to `in_progress` before subscribing or sending, and never replays a claimed or interrupted Turn. Extra inputs require native receipts. The terminal outcome and the native Session ID commit together under the admission lock, and unapplied messages prevent a successful completion. Credentials are resolved separately from the immutable non-secret snapshot. A peer that lacks a required capability is rejected before the claim, and a failed strict resume never starts unrelated history. Native continuation needs the device's persisted engine files; a native Session ID alone cannot restore deleted history.

Core replaces a Turn's usage with each complete valid token breakdown and keeps the last committed measurement when execution is interrupted. It never infers tokens from context occupancy or costs and never parses native raw payloads. When `done` also carries usage that was already reported, the counters are not added again.
Core replaces a Turn's usage with each complete valid token breakdown, which `sessions.MeasuredUsage` decides, and keeps the last committed measurement when execution is interrupted. It never infers tokens from context occupancy or costs and never parses native raw payloads. When `done` also carries usage that was already reported, the counters are not added again.

Subagent observations use the common types in `internal/agentdaemon/proto/subagents.go`. Core assigns public IDs and projects them under the Session lock and the leased execution journal; native names, history parsing and outcome proof stay in the adapters. Child Turns have their own table, separate from Core's queue, and Session Turn reads and the Session event stream carry root work only. The migration-defined `public_execution_turns` view has no public reader; do not reintroduce mixed Session Turn pages. Repeated effects are idempotent, root output freezes first, and child Items are delivered before their terminal Turn snapshot.

Expand All @@ -151,15 +153,15 @@ Neutral tool and message observations go into the journal before Core projects p

Execution observations are written to tenant-scoped `turn_events` in ordered, idempotent batches before they back recovery or publication, with daemon payloads intact. Flush at least every 100 ms while consuming events and before terminal persistence; the terminal outcome, its journal entry and native continuity commit together. Journal limits are 512 KiB per payload, 1 MiB per batch, 65,536 observations and 32 MiB per Turn, with one extra entry reserved for the terminal outcome. Never infer a successful completion after a persistence error or stream overflow.

Live Session SSE reads `session_events` committed with the corresponding input, Item or lifecycle change under the Session lock, as immutable transition snapshots. The notification buffer keeps at most 256 events and 64 MiB per Session after each transaction, and read batches at most 32 events or 1 MiB, each keeping a single oversized event. GET polls committed events every 100 ms from the committed high-water mark; a missing sequence position ends the stream with a safe error. Socket writes have a five-second deadline and hold no database connection. Rebuilding historical indexes emits no live events.
Live Session SSE reads `session_events` committed with the corresponding input, Item or lifecycle change under the Session lock, as immutable transition snapshots. The notification buffer keeps at most 256 events and 64 MiB per Session after each transaction (the `sessions` retention bounds, which `sessionpg` applies), and read batches at most 32 events or 1 MiB, each keeping a single oversized event. GET polls committed events every 100 ms from the committed high-water mark; a missing sequence position ends the stream with a safe error. Socket writes have a five-second deadline and hold no database connection. Rebuilding historical indexes emits no live events.

A streaming Session creation reuses atomic input admission and the live event loop. The creation upsert returns its event cursor under the Session lock, before the initial inputs; never replace it with a post-commit cursor lookup. The settled marker and pending-input flag the stream uses to end are store-internal, never wire fields. The stream rereads the Session projection after it sends a Session status event and otherwise at most once a second.
A streaming Session creation reuses atomic input admission and the live event loop. The creation upsert returns its event cursor under the Session lock, before the initial inputs; never replace it with a post-commit cursor lookup. The settled marker on `sessions.SessionChange` and the pending-input flag the stream uses to end are internal, never wire fields. The stream rereads the Session projection after it sends a Session status event and otherwise at most once a second.

## Items

Public Items read a projection updated in the same Session transaction as admitted messages and journal batches. Item IDs derive from the Turn and source identity, and the first-observation timestamp and tie breakers never change when content or status does. Each new Item's Session position is allocated under the Session lock, preserving observation order for equal timestamps, and each Turn allocates its own zero-based `output_index`, which inputs do not consume; updates and retries keep both.

Item merging never mutates the incoming observation or the previous snapshot: public text delta events read the original fragment after merging, while the Item keeps the accumulated text, and the content slice is copied before its text pointer is replaced. A first observation without its own fragment carries its unchanged text in one delta. Wire-only explicit nulls come from response marshalling, while stored Item payloads keep their original encoding through `Item.MarshalStored`, so replayed child Items compare equal. Structured tool JSON is kept without float conversion, and an unfinished call never becomes a successful result. Function results are Session input Items: they emit `item.added` with a null `output_index` and never `item.done`, whose upstream union allows only agent output, and their public output and error come from the saved submission.
Item merging never mutates the incoming observation or the previous snapshot: public text delta events read the original fragment after merging, while the Item keeps the accumulated text, and the content slice is copied before its text pointer is replaced. A first observation without its own fragment carries its unchanged text in one delta. Wire-only explicit nulls come from response marshalling, while stored Item payloads keep their original encoding through `Item.MarshalStored`, so replayed child Items compare equal. Structured tool JSON is kept without float conversion, and an unfinished call never becomes a successful result. When a Turn ends, `sessions.EndTurn` makes its unfinished Items incomplete with their partial content and reports them in Session position order before the Turn's event and its settled Session activity; `sessionpg` gives them one shared settlement time. Function results are Session input Items: they emit `item.added` with a null `output_index` and never `item.done`, whose upstream union allows only agent output, and their public output and error come from the saved submission.

## Worker ownership

Expand Down
3 changes: 2 additions & 1 deletion services/core/internal/api/admin_resources_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import (
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/adminaudit"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/agents"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/identity"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/store"
)

Expand Down Expand Up @@ -123,7 +124,7 @@ func (s *summaryFixture) ReadAdminSummary(_ context.Context, tenant string, filt
for i, usage := range []json.RawMessage{nil, json.RawMessage(`{"input_tokens":3,"output_tokens":5,"total_tokens":8,"input_tokens_details":{"cached_tokens":2},"output_tokens_details":{"reasoning_tokens":1}}`)} {
session := store.Session{ID: "session", TenantID: tenant, Configuration: json.RawMessage(`{"agent":{"id":"agent","model":"model","tools":[]},"environment":{"type":"none"}}`), CreatedAt: time.Unix(100+int64(i), 0), Usage: usage}
if i == 0 {
session.LastTurn = &store.Turn{Status: store.TurnInProgress, CreatedAt: time.Unix(110, 0)}
session.LastTurn = &sessions.Turn{Status: sessions.TurnInProgress, CreatedAt: time.Unix(110, 0)}
}
if err := visit(session, nil); err != nil {
return store.AdminAssetCounts{}, err
Expand Down
9 changes: 5 additions & 4 deletions services/core/internal/api/fakes_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ import (
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimehistory"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/store"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/writeaudit"
)
Expand Down Expand Up @@ -739,7 +740,7 @@ func (f *fakeSessionArchive) ArchiveManagedSession(a0 context.Context, a1 string
type fakeSessionEvents struct {
t testing.TB
sessionEventCursor func(context.Context, string, string) (int64, error)
listSessionEvents func(context.Context, string, string, int64) ([]store.SessionChange, error)
listSessionEvents func(context.Context, string, string, int64) ([]sessions.SessionChange, error)
sessionStreamSnapshot func(context.Context, string, string) (store.Session, int64, error)
}

Expand All @@ -750,7 +751,7 @@ func (f *fakeSessionEvents) SessionEventCursor(a0 context.Context, a1 string, a2
return f.sessionEventCursor(a0, a1, a2)
}

func (f *fakeSessionEvents) ListSessionEvents(a0 context.Context, a1 string, a2 string, a3 int64) ([]store.SessionChange, error) {
func (f *fakeSessionEvents) ListSessionEvents(a0 context.Context, a1 string, a2 string, a3 int64) ([]sessions.SessionChange, error) {
if f.listSessionEvents == nil {
unexpectedCall(f.t, "ListSessionEvents")
}
Expand All @@ -766,12 +767,12 @@ func (f *fakeSessionEvents) SessionStreamSnapshot(a0 context.Context, a1 string,

type fakeSessionHistory struct {
t testing.TB
getTurn func(context.Context, string, string, string) (store.Turn, error)
getTurn func(context.Context, string, string, string) (sessions.Turn, error)
listTurns func(context.Context, string, string, string, int, bool) (store.TurnPage, error)
listItems func(context.Context, string, string, string, int, bool) (store.ItemPage, error)
}

func (f *fakeSessionHistory) GetTurn(a0 context.Context, a1 string, a2 string, a3 string) (store.Turn, error) {
func (f *fakeSessionHistory) GetTurn(a0 context.Context, a1 string, a2 string, a3 string) (sessions.Turn, error) {
if f.getTurn == nil {
unexpectedCall(f.t, "GetTurn")
}
Expand Down
5 changes: 3 additions & 2 deletions services/core/internal/api/function_state_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (
"time"

v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/store"
)

Expand All @@ -16,7 +17,7 @@ func TestFunctionStateEventsUseTheirOwnSnapshot(t *testing.T) {
}
for _, arguments := range []string{`{"n":9007199254740993}`, `null`, `[1,"value"]`, `"value"`, `false`} {
action := v1.FunctionCallAction{Type: "function_call", CallID: "call", Name: "lookup", TurnID: "turn", Arguments: json.RawMessage(arguments)}
change := store.SessionChange{Event: v1.SessionEvent{Type: "agent.session.requires_action", EventID: "event", SessionID: "session"}, Turn: &store.Turn{ID: "turn", Status: store.TurnWaiting}, RequiredActions: []v1.FunctionCallAction{action}}
change := sessions.SessionChange{Event: v1.SessionEvent{Type: "agent.session.requires_action", EventID: "event", SessionID: "session"}, Turn: &sessions.Turn{ID: "turn", Status: sessions.TurnWaiting}, RequiredActions: []v1.FunctionCallAction{action}}
event, err := streamResponse(session, change, "")
if err != nil || event.Session.Status != "requires_action" {
t.Fatal(event, err)
Expand All @@ -38,7 +39,7 @@ func TestFunctionStateEventsUseTheirOwnSnapshot(t *testing.T) {
t.Fatal("unexpected public fields", string(raw))
}
change.RequiredActions = nil
change.Turn.Status = store.TurnInProgress
change.Turn.Status = sessions.TurnInProgress
change.Event.Type = "agent.session.in_progress"
event, err = streamResponse(session, change, "")
if err != nil || event.Session.Status != "in_progress" || event.Session.RequiredActions == nil || len(event.Session.RequiredActions) != 0 {
Expand Down
Loading
Loading