diff --git a/contracts/agents-api/README.md b/contracts/agents-api/README.md index c331ddb01..5049296c6 100644 --- a/contracts/agents-api/README.md +++ b/contracts/agents-api/README.md @@ -98,6 +98,16 @@ paths start at `/vaults`, not `/agents/vaults`. | vaults | create, retrieve, list, delete | Create/retrieve/list/delete with independent tenant persistence, stored status filtering, atomic Credential cascade and frozen Session attachments; archive semantics and full hosted lifecycle parity remain missing | | vaults.credentials | create, retrieve, update, list, delete | Static-bearer create/retrieve/list/token replacement/deletion with scoped encrypted storage; Session attachment and exact-URL HTTPS MCP binding; OAuth, archive semantics and full hosted lifecycle parity remain missing | +## Core extension inventory + +The operations below are implemented public Core extensions. They are excluded +from the 42-operation upstream inventory and must not be counted as OpenAI Agents +compatibility. + +| Extension | Operations | Current coverage | +| --- | --- | --- | +| Runtime observations | `GET /v1/agents/runtime-observations`; `GET /v1/agents/sessions/{session_id}/runtime-observation` | Current, read-only, tenant-scoped Session contexts with stable Session-keyset pagination, bounded concurrent sampling, Docker metrics, explicit unsupported/unavailable states, strict `packages/agents-client` projection, and no lifecycle mutation. Kubernetes, E2B, self-hosted telemetry, history, CPU-rate derivation, and automatic idle policy remain unimplemented. See [Runtime observation API](runtime-observability-api.md). | + For each resource, verify the referenced request/response unions and observable behavior, not just the route. Non-text initial input, configuration options, text/image content, function results, environment variants, full Item/SSE diff --git a/contracts/agents-api/openapi.yaml b/contracts/agents-api/openapi.yaml index 3578bf4ff..90a206045 100644 --- a/contracts/agents-api/openapi.yaml +++ b/contracts/agents-api/openapi.yaml @@ -868,6 +868,180 @@ definitions: required: - type type: object + v1.RuntimeCPUObservation: + properties: + capacity_cores: + minimum: 5e-324 + type: number + x-nullable: true + usage_cores: + minimum: 0 + type: number + x-nullable: true + usage_seconds_total: + minimum: 0 + type: number + x-nullable: true + utilization_ratio: + minimum: 0 + type: number + x-nullable: true + required: + - capacity_cores + - usage_cores + - usage_seconds_total + - utilization_ratio + type: object + v1.RuntimeInstance: + properties: + allocation_id: + format: uuid + type: string + x-nullable: true + connection_generation: + format: uuid + type: string + x-nullable: true + device_id: + format: uuid + type: string + x-nullable: true + kind: + enum: + - managed_allocation + - self_hosted_connection + - none + type: string + required: + - allocation_id + - connection_generation + - device_id + - kind + type: object + v1.RuntimeMemoryObservation: + properties: + limit_bytes: + minimum: 1 + type: integer + x-nullable: true + usage_bytes: + minimum: 0 + type: integer + x-nullable: true + required: + - limit_bytes + - usage_bytes + type: object + v1.RuntimeObservation: + properties: + allocation_created_at: + minimum: 0 + type: integer + x-nullable: true + cpu: + allOf: + - $ref: '#/definitions/v1.RuntimeCPUObservation' + x-nullable: true + environment_id: + format: uuid + type: string + x-nullable: true + id: + format: uuid + type: string + instance: + $ref: '#/definitions/v1.RuntimeInstance' + memory: + allOf: + - $ref: '#/definitions/v1.RuntimeMemoryObservation' + x-nullable: true + mode: + enum: + - none + - self_hosted + - openai_hosted + type: string + object: + enum: + - agent.runtime_observation + type: string + observed_at: + minimum: 0 + type: integer + x-nullable: true + provider_type: + type: string + x-nullable: true + reason: + enum: + - runtime_mode_not_observable + - allocation_pending + - runtime_not_running + - source_not_configured + - sample_timeout + - sample_unavailable + type: string + x-nullable: true + resolved_at: + minimum: 0 + type: integer + session_id: + format: uuid + type: string + started_at: + minimum: 0 + type: integer + x-nullable: true + status: + enum: + - observed + - unsupported + - unavailable + type: string + required: + - allocation_created_at + - cpu + - environment_id + - id + - instance + - memory + - mode + - object + - observed_at + - provider_type + - reason + - resolved_at + - session_id + - started_at + - status + type: object + v1.RuntimeObservationList: + properties: + data: + items: + $ref: '#/definitions/v1.RuntimeObservation' + type: array + first_id: + format: uuid + type: string + x-nullable: true + has_more: + type: boolean + last_id: + format: uuid + type: string + x-nullable: true + object: + enum: + - list + type: string + required: + - data + - first_id + - has_more + - last_id + - object + type: object v1.SavedAgent: properties: created_at: @@ -2581,6 +2755,68 @@ paths: summary: Update an Environment Template tags: - Environment Templates + /agents/runtime-observations: + get: + description: Core extension listing one current Runtime context per tenant-owned + Session in Session creation order. Each row has an independent resolved_at + and optional provider observed_at; the page is not an atomic telemetry snapshot. + parameters: + - description: agents=v1 + in: header + name: OpenAI-Beta + required: true + type: string + - description: Last observation ID from the previous page + in: query + name: after + type: string + - default: 20 + description: Page size + in: query + maximum: 100 + minimum: 1 + name: limit + type: integer + - default: desc + description: Session creation order + enum: + - asc + - desc + in: query + name: order + type: string + produces: + - application/json + responses: + "200": + description: OK + schema: + $ref: '#/definitions/v1.RuntimeObservationList' + "400": + description: Bad Request + schema: + $ref: '#/definitions/v1.ErrorResponse' + "401": + description: Unauthorized + schema: + $ref: '#/definitions/v1.ErrorResponse' + "404": + description: Not Found + schema: + $ref: '#/definitions/v1.ErrorResponse' + "500": + description: Internal Server Error + schema: + $ref: '#/definitions/v1.ErrorResponse' + "503": + description: Service Unavailable + schema: + $ref: '#/definitions/v1.ErrorResponse' + security: + - BearerAuth: [] + summary: List current Runtime observations + tags: + - Runtime observations /agents/sessions: get: description: Cursor and results are scoped to the authenticated execution tenant. @@ -3348,6 +3584,53 @@ paths: summary: List persisted execution Items tags: - Items + /agents/sessions/{session_id}/runtime-observation: + get: + description: Core extension returning one tenant-scoped, read-only current Runtime + observation. It never provisions, renews, restarts, pauses or stops compute. + parameters: + - description: agents=v1 + in: header + name: OpenAI-Beta + required: true + type: string + - description: Session ID + in: path + name: session_id + required: true + type: string + produces: + - application/json + responses: + "200": + description: OK + schema: + $ref: '#/definitions/v1.RuntimeObservation' + "400": + description: Bad Request + schema: + $ref: '#/definitions/v1.ErrorResponse' + "401": + description: Unauthorized + schema: + $ref: '#/definitions/v1.ErrorResponse' + "404": + description: Not Found + schema: + $ref: '#/definitions/v1.ErrorResponse' + "500": + description: Internal Server Error + schema: + $ref: '#/definitions/v1.ErrorResponse' + "503": + description: Service Unavailable + schema: + $ref: '#/definitions/v1.ErrorResponse' + security: + - BearerAuth: [] + summary: Retrieve a Session Runtime observation + tags: + - Runtime observations /agents/sessions/{session_id}/subagents: get: description: Includes nested and closed Subagents. Cursors belong to the same diff --git a/contracts/agents-api/runtime-observability-api.md b/contracts/agents-api/runtime-observability-api.md index 4e607a9d7..34394ffd3 100644 --- a/contracts/agents-api/runtime-observability-api.md +++ b/contracts/agents-api/runtime-observability-api.md @@ -1,7 +1,9 @@ # Runtime observation API proposal -Status: review proposal. These routes are not implemented and are not yet present -in `openapi.yaml`. +Status: Phase 2 implemented. The current-snapshot routes, strict +`packages/agents-client` projection, and generated `openapi.yaml` contract are +implemented. Core Web integration, historical queries, and lifecycle controls +remain outside this phase. This is an Agents Core extension, not an upstream OpenAI Agents resource. The implementation must record that status in the coverage ledger and generated @@ -40,13 +42,13 @@ before publishing a new Dashboard snapshot. "id": "6c77d3a2-71d6-4ed5-884f-687aecda02a3", "object": "agent.runtime_observation", "session_id": "6c77d3a2-71d6-4ed5-884f-687aecda02a3", - "environment_id": "env_...", + "environment_id": "6c02fb71-5fa8-4298-93e8-57c6625a3fc2", "mode": "openai_hosted", "provider_type": "docker", "instance": { "kind": "managed_allocation", - "allocation_id": "alloc_...", - "device_id": "device_...", + "allocation_id": "d23ab94e-e40b-45bd-93a2-444f1f74642b", + "device_id": "2e434f4f-76aa-4e54-a707-4757036d90ef", "connection_generation": null }, "status": "observed", @@ -58,8 +60,8 @@ before publishing a new Dashboard snapshot. "cpu": { "usage_seconds_total": 482.75, "capacity_cores": 2.0, - "usage_cores": 1.42, - "utilization_ratio": 0.71 + "usage_cores": null, + "utilization_ratio": null }, "memory": { "usage_bytes": 805306368, @@ -181,7 +183,6 @@ Use the existing Agents API error envelope. | 400 | `invalid_request` | Empty or invalid limits, order, or malformed cursor. | | 401 | `authentication_error` | Missing or invalid API authentication. | | 404 | `not_found` | Missing or foreign Session/cursor, indistinguishably. | -| 429 | `rate_limit_exceeded` | Runtime sampling read budget exceeded. | | 500 | `internal_error` | Integrity, ownership, or invalid provider evidence. | | 503 | `execution_unavailable` | Required Runtime observation service is not configured. | @@ -190,14 +191,15 @@ Errors never include provider raw responses or credentials. ## Freshness and caching - Return `Cache-Control: no-store`. -- An internal cache may coalesce reads for at most five seconds. +- The Phase 2 implementation performs bounded direct reads and has no observation + cache. A later internal cache may coalesce reads for at most five seconds. - `observed_at` is authoritative for freshness; HTTP response time is not. - Clients mark samples stale according to their own explicit threshold. - `ETag` is not proposed because observations change independently. ## Client contract -`packages/agents-client` should expose: +`packages/agents-client` exposes: ```ts type RuntimeObservationStatus = "observed" | "unsupported" | "unavailable"; @@ -209,38 +211,13 @@ type RuntimeObservationReason = | "sample_timeout" | "sample_unavailable"; -interface RuntimeObservation { - id: string; - object: "agent.runtime_observation"; - session_id: string; - environment_id: string | null; - mode: "none" | "self_hosted" | "openai_hosted"; - provider_type: string | null; - instance: { - kind: "managed_allocation" | "self_hosted_connection" | "none"; - allocation_id: string | null; - device_id: string | null; - connection_generation: string | null; - }; - status: RuntimeObservationStatus; - reason: RuntimeObservationReason | null; - allocation_created_at: number | null; - resolved_at: number; - observed_at: number | null; - started_at: number | null; - cpu: { - usage_seconds_total: number | null; - capacity_cores: number | null; - usage_cores: number | null; - utilization_ratio: number | null; - } | null; - memory: { - usage_bytes: number | null; - limit_bytes: number | null; - } | null; -} +type RuntimeObservation = + | RuntimeObservedObservation + | RuntimeUnavailableObservation + | RuntimeNoneObservation + | RuntimeSelfHostedObservation; -interface RuntimeObservationPage { +interface RuntimeObservationList { object: "list"; data: RuntimeObservation[]; has_more: boolean; @@ -248,20 +225,33 @@ interface RuntimeObservationPage { last_id: string | null; } -interface RuntimeObservationClient { - list(options?: { +interface AgentCore { + listRuntimeObservations(options?: { after?: string; limit?: number; order?: "asc" | "desc"; - }): Promise; + }): Promise; - retrieveForSession(sessionId: string): Promise; + retrieveRuntimeObservation(sessionId: string): Promise; } ``` +These exported variants discriminate on `status` and `mode`; their instance, +reason, timestamps, CPU, and memory fields narrow accordingly. The exact variant +definitions live in `packages/agents-client/src/types.ts` and mirror the status +and reason matrix above. + The client validates every required field, enum, nullability rule, timestamp, and -finite number. Unknown additive fields are ignored. Malformed data rejects the -whole page; Web does not publish a partial snapshot. +finite number. The current pinned contract rejects unknown additive fields so an +unreviewed server expansion cannot silently cross the browser boundary. Malformed +data rejects the whole page; Web does not publish a partial snapshot. + +The generated OpenAPI 2 schema records field-level required/nullability rules, +UUID formats, reason enums, and numeric minima. OpenAPI 2 +cannot encode the complete cross-field discriminated union. The matrix above is +normative for wire consumers; the server projection and strict TypeScript +projector enforce it, and the exported TypeScript type prevents invalid +status/mode combinations in typed consumers. Web also applies a configured whole-refresh budget. If `has_more` remains true when that budget is exhausted, it retains the prior complete snapshot and marks diff --git a/contracts/agents-api/runtime-observability-design.md b/contracts/agents-api/runtime-observability-design.md index 70ca2fe63..dd9dfa16d 100644 --- a/contracts/agents-api/runtime-observability-design.md +++ b/contracts/agents-api/runtime-observability-design.md @@ -1,8 +1,8 @@ # Runtime observability and Dashboard design -Status: review proposal. Phase 1 provider abstraction and Docker sampling are -implemented; the public API, Web integration, history backend, additional -providers, and lifecycle automation described below are not implemented. +Status: Phase 1 provider abstraction/Docker sampling and Phase 2 current-snapshot +API/client contract are implemented. Core Web integration, history backend, +additional providers, and lifecycle automation described below are not implemented. ## 1. Problem statement @@ -156,9 +156,11 @@ observation times. A single sample cannot truthfully supply CPU percentage. The API projection may additionally expose `usage_cores` and `utilization_ratio` only when the service has two ordered samples for the same Runtime incarnation. -The process-local observation cache keeps the previous cumulative value for this -calculation. Its loss makes the derived fields temporarily null; it never changes -the cumulative source measurement or lifecycle state. +A future process-local observation cache may keep the previous cumulative value +for this calculation. Phase 2 intentionally leaves both derived fields null +because it has only one provider sample per request. Cache loss must make the +derived fields temporarily null; it must never change the cumulative source +measurement or lifecycle state. ## 8. Duration semantics @@ -324,6 +326,8 @@ transitions. That migration cannot read a monitoring backend as authority. ### Phase 1: provider-neutral foundation +Implemented in the Docker observability foundation. + - `runtimeobs` identity, resolver, source, sample, and service. - Managed Docker Inspect/Stats source. - CPU, memory, and current compute start time. @@ -331,6 +335,10 @@ transitions. That migration cannot read a monitoring backend as authority. ### Phase 2: current snapshot API +Implemented by the Runtime Observation extension routes and +`packages/agents-client`. The generated OpenAPI contract records the extension; +this does not add an upstream OpenAI operation. + - Add extension types under `contracts/agents-api/v1`. - Add collection and Session-scoped handlers. - Add `packages/agents-client` methods and raw HTTP/client coverage. @@ -376,13 +384,12 @@ transitions. That migration cannot read a monitoring backend as authority. - No lifecycle action is reachable from the first Dashboard. - No new database table is required for current snapshots or history export. -## 17. Review decisions required +## 17. Recorded design decisions -1. Accept the proposed API as a documented Core extension rather than an upstream - OpenAI resource. -2. Accept current snapshots without atomic cross-row time semantics; every row +1. The API is a documented Core extension rather than an upstream OpenAI resource. +2. Current snapshots have no atomic cross-row time semantics; every row exposes its own `observed_at`. -3. Confirm that history is optional and external, not a PostgreSQL sample table. -4. Confirm that the first Web release has no lifecycle controls. -5. Choose whether `self_hosted` remains visibly unsupported until authenticated +3. History is optional and external, not a PostgreSQL sample table. +4. The first Web release has no lifecycle controls. +5. `self_hosted` remains visibly unsupported until authenticated, generation-fenced telemetry is qualified. diff --git a/contracts/agents-api/v1/runtime_observations.go b/contracts/agents-api/v1/runtime_observations.go new file mode 100644 index 000000000..d82d8d514 --- /dev/null +++ b/contracts/agents-api/v1/runtime_observations.go @@ -0,0 +1,46 @@ +package v1 + +type RuntimeObservation struct { + ID string `json:"id" binding:"required" format:"uuid"` + Object string `json:"object" enums:"agent.runtime_observation" binding:"required"` + SessionID string `json:"session_id" binding:"required" format:"uuid"` + EnvironmentID *string `json:"environment_id" extensions:"x-nullable" binding:"required" format:"uuid"` + Mode string `json:"mode" enums:"none,self_hosted,openai_hosted" binding:"required"` + ProviderType *string `json:"provider_type" extensions:"x-nullable" binding:"required" pattern:"^[a-z][a-z0-9_]{0,31}$"` + Instance RuntimeInstance `json:"instance" binding:"required"` + Status string `json:"status" enums:"observed,unsupported,unavailable" binding:"required"` + Reason *string `json:"reason" extensions:"x-nullable" binding:"required" enums:"runtime_mode_not_observable,allocation_pending,runtime_not_running,source_not_configured,sample_timeout,sample_unavailable"` + AllocationCreatedAt *int64 `json:"allocation_created_at" extensions:"x-nullable" binding:"required" minimum:"0"` + ResolvedAt int64 `json:"resolved_at" binding:"required" minimum:"0"` + ObservedAt *int64 `json:"observed_at" extensions:"x-nullable" binding:"required" minimum:"0"` + StartedAt *int64 `json:"started_at" extensions:"x-nullable" binding:"required" minimum:"0"` + CPU *RuntimeCPUObservation `json:"cpu" extensions:"x-nullable" binding:"required"` + Memory *RuntimeMemoryObservation `json:"memory" extensions:"x-nullable" binding:"required"` +} + +type RuntimeInstance struct { + Kind string `json:"kind" enums:"managed_allocation,self_hosted_connection,none" binding:"required"` + AllocationID *string `json:"allocation_id" extensions:"x-nullable" binding:"required" format:"uuid"` + DeviceID *string `json:"device_id" extensions:"x-nullable" binding:"required" format:"uuid"` + ConnectionGeneration *string `json:"connection_generation" extensions:"x-nullable" binding:"required" format:"uuid"` +} + +type RuntimeCPUObservation struct { + UsageSecondsTotal *float64 `json:"usage_seconds_total" extensions:"x-nullable" binding:"required" minimum:"0"` + CapacityCores *float64 `json:"capacity_cores" extensions:"x-nullable" binding:"required" minimum:"5e-324"` + UsageCores *float64 `json:"usage_cores" extensions:"x-nullable" binding:"required" minimum:"0"` + UtilizationRatio *float64 `json:"utilization_ratio" extensions:"x-nullable" binding:"required" minimum:"0"` +} + +type RuntimeMemoryObservation struct { + UsageBytes *uint64 `json:"usage_bytes" extensions:"x-nullable" binding:"required" minimum:"0"` + LimitBytes *uint64 `json:"limit_bytes" extensions:"x-nullable" binding:"required" minimum:"1"` +} + +type RuntimeObservationList struct { + Object string `json:"object" enums:"list" binding:"required"` + Data []RuntimeObservation `json:"data" binding:"required"` + HasMore bool `json:"has_more" binding:"required"` + FirstID *string `json:"first_id" extensions:"x-nullable" binding:"required" format:"uuid"` + LastID *string `json:"last_id" extensions:"x-nullable" binding:"required" format:"uuid"` +} diff --git a/packages/agents-client/src/client.test.ts b/packages/agents-client/src/client.test.ts index 13c0a237d..facd5e497 100644 --- a/packages/agents-client/src/client.test.ts +++ b/packages/agents-client/src/client.test.ts @@ -103,6 +103,42 @@ function messageItem(overrides: Record = {}): Record = {}): Record { + return { + id: runtimeSessionId, + object: "agent.runtime_observation", + session_id: runtimeSessionId, + environment_id: runtimeEnvironmentId, + mode: "openai_hosted", + provider_type: "docker", + instance: { + kind: "managed_allocation", + allocation_id: runtimeAllocationId, + device_id: runtimeDeviceId, + connection_generation: null, + }, + status: "observed", + reason: null, + allocation_created_at: 10, + resolved_at: 30, + observed_at: 20, + started_at: 10, + cpu: { + usage_seconds_total: 0, + capacity_cores: 2, + usage_cores: null, + utilization_ratio: null, + }, + memory: { usage_bytes: 0, limit_bytes: 1024 }, + ...overrides, + }; +} + describe("OpenAIAgentsClient", () => { afterEach(() => vi.unstubAllGlobals()); @@ -2385,4 +2421,121 @@ describe("OpenAIAgentsClient", () => { ).rejects.toThrow("createSession only supports the JSON response"); expect(calls).toHaveLength(0); }); + + it("retrieves a Runtime observation, preserves observed zeroes, and encodes the Session ID", async () => { + const calls: FetchCall[] = []; + const client = new OpenAIAgentsClient({ + baseUrl: "https://core.example/v1", + fetch: recordingFetch(jsonResponse(runtimeObservation()), calls), + }); + + await expect(client.retrieveRuntimeObservation(runtimeSessionId)).resolves.toMatchObject({ + id: runtimeSessionId, + cpu: { usage_seconds_total: 0 }, + memory: { usage_bytes: 0 }, + }); + expect(String(calls[0]?.input)).toBe( + `https://core.example/v1/agents/sessions/${runtimeSessionId}/runtime-observation`, + ); + }); + + it("lists Runtime observations with stable pagination metadata and query serialization", async () => { + const calls: FetchCall[] = []; + const body = { + object: "list", + data: [runtimeObservation()], + has_more: true, + first_id: runtimeSessionId, + last_id: runtimeSessionId, + }; + const client = new OpenAIAgentsClient({ + baseUrl: "https://core.example/v1/", + fetch: recordingFetch(jsonResponse(body), calls), + }); + + await expect(client.listRuntimeObservations({ + after: runtimeSessionId, limit: 1, order: "asc", + })).resolves.toMatchObject(body); + expect(String(calls[0]?.input)).toBe( + `https://core.example/v1/agents/runtime-observations?after=${runtimeSessionId}&limit=1&order=asc`, + ); + }); + + it("accepts an unsupported none-mode Runtime observation with explicit nulls", async () => { + const value = runtimeObservation({ + environment_id: null, + mode: "none", + provider_type: null, + instance: { kind: "none", allocation_id: null, device_id: null, connection_generation: null }, + status: "unsupported", + reason: "runtime_mode_not_observable", + allocation_created_at: null, + observed_at: null, + started_at: null, + cpu: null, + memory: null, + }); + const client = new OpenAIAgentsClient({ fetch: recordingFetch(jsonResponse(value), []) }); + await expect(client.retrieveRuntimeObservation(runtimeSessionId)).resolves.toMatchObject(value); + }); + + it.each([ + ["unknown field", () => ({ ...runtimeObservation(), provider_native_id: "hidden" })], + ["foreign Session", () => ({ ...runtimeObservation(), session_id: "55555555-5555-4555-8555-555555555555" })], + ["invalid status/reason", () => ({ ...runtimeObservation(), status: "observed", reason: "sample_timeout" })], + ["invalid mode/instance", () => ({ ...runtimeObservation(), mode: "none" })], + ["negative CPU", () => ({ ...runtimeObservation(), cpu: { + usage_seconds_total: -1, capacity_cores: 2, usage_cores: null, utilization_ratio: null, + } })], + ["non-numeric CPU", () => ({ ...runtimeObservation(), cpu: { + usage_seconds_total: "NaN", capacity_cores: 2, usage_cores: null, utilization_ratio: null, + } })], + ["zero CPU capacity", () => ({ ...runtimeObservation(), cpu: { + usage_seconds_total: 1, capacity_cores: 0, usage_cores: null, utilization_ratio: null, + } })], + ["unsafe memory", () => ({ ...runtimeObservation(), memory: { + usage_bytes: Number.MAX_SAFE_INTEGER + 1, limit_bytes: 1024, + } })], + ["zero memory limit", () => ({ ...runtimeObservation(), memory: { + usage_bytes: 1, limit_bytes: 0, + } })], + ])("rejects a Runtime observation with %s", async (_label, build) => { + const client = new OpenAIAgentsClient({ fetch: recordingFetch(jsonResponse(build()), []) }); + await expect(client.retrieveRuntimeObservation(runtimeSessionId)).rejects.toMatchObject({ + status: 502, + code: "invalid_runtime_observation", + }); + }); + + it.each([ + ["mismatched first_id", { + object: "list", data: [runtimeObservation()], has_more: false, + first_id: runtimeEnvironmentId, last_id: runtimeSessionId, + }], + ["duplicate IDs", { + object: "list", data: [runtimeObservation(), runtimeObservation()], has_more: false, + first_id: runtimeSessionId, last_id: runtimeSessionId, + }], + ["empty continuation", { + object: "list", data: [], has_more: true, first_id: null, last_id: null, + }], + ])("rejects a Runtime observation list with %s", async (_label, body) => { + const client = new OpenAIAgentsClient({ fetch: recordingFetch(jsonResponse(body), []) }); + await expect(client.listRuntimeObservations()).rejects.toMatchObject({ + status: 502, + code: "invalid_runtime_observation", + }); + }); + + it.each([ + { after: "not-a-uuid" }, + { limit: 0 }, + { limit: 101 }, + { order: "sideways" }, + ])("rejects invalid Runtime observation pagination before fetch", async (options) => { + const calls: FetchCall[] = []; + const client = new OpenAIAgentsClient({ fetch: recordingFetch(jsonResponse({}), calls) }); + await expect(client.listRuntimeObservations(options as never)).rejects.toThrow(TypeError); + expect(calls).toHaveLength(0); + }); }); diff --git a/packages/agents-client/src/client.ts b/packages/agents-client/src/client.ts index c2f32630b..dc4393f57 100644 --- a/packages/agents-client/src/client.ts +++ b/packages/agents-client/src/client.ts @@ -43,6 +43,8 @@ import type { StreamError, UpdateAgentInput, ReplaceVaultCredentialTokenInput, + RuntimeObservation, + RuntimeObservationList, Vault, VaultCredential, VaultCredentialDeleted, @@ -234,6 +236,18 @@ const knownItemTypes = new Set([ ]); const itemStatuses = new Set(["in_progress", "completed", "failed", "incomplete"]); const turnStatuses = new Set(["queued", "in_progress", "waiting", "completed", "failed", "cancelled"]); +const runtimeObservationFields = new Set([ + "id", "object", "session_id", "environment_id", "mode", "provider_type", "instance", "status", "reason", + "allocation_created_at", "resolved_at", "observed_at", "started_at", "cpu", "memory", +]); +const runtimeInstanceFields = new Set(["kind", "allocation_id", "device_id", "connection_generation"]); +const runtimeCPUFields = new Set(["usage_seconds_total", "capacity_cores", "usage_cores", "utilization_ratio"]); +const runtimeMemoryFields = new Set(["usage_bytes", "limit_bytes"]); +const runtimeObservationReasons = new Set([ + "runtime_mode_not_observable", "allocation_pending", "runtime_not_running", + "source_not_configured", "sample_timeout", "sample_unavailable", +]); +const runtimeProviderTypePattern = /^[a-z][a-z0-9_]{0,31}$/; function exactFields(value: Record, fields: Set): boolean { const keys = Object.keys(value); @@ -995,6 +1009,164 @@ function projectAgentSession( return session; } +function invalidRuntimeObservation(message = "Agent Core returned an invalid Runtime observation."): never { + throw new AgentCoreError(message, 502, "invalid_runtime_observation"); +} + +function nullableRuntimeNumber(value: unknown): number | null { + if (value === null) return null; + if (typeof value !== "number" || !Number.isFinite(value) || value < 0) { + return invalidRuntimeObservation(); + } + return value; +} + +function nullableRuntimeInteger(value: unknown): number | null { + const projected = nullableRuntimeNumber(value); + if (projected !== null && !Number.isSafeInteger(projected)) return invalidRuntimeObservation(); + return projected; +} + +function projectRuntimeObservation(value: unknown, expectedSessionId?: string): RuntimeObservation { + if (!isRecord(value) || !exactFields(value, runtimeObservationFields)) { + return invalidRuntimeObservation(); + } + const id = canonicalUuid(value.id); + const sessionId = canonicalUuid(value.session_id); + const environmentId = value.environment_id === null ? null : canonicalUuid(value.environment_id); + if ( + id === null || sessionId === null || id !== sessionId || + (expectedSessionId !== undefined && !sameUuid(sessionId, expectedSessionId)) || + value.object !== "agent.runtime_observation" || + (value.mode !== "none" && value.mode !== "self_hosted" && value.mode !== "openai_hosted") || + !(value.provider_type === null || ( + typeof value.provider_type === "string" && runtimeProviderTypePattern.test(value.provider_type) + )) || + !isRecord(value.instance) || !exactFields(value.instance, runtimeInstanceFields) || + (value.status !== "observed" && value.status !== "unsupported" && value.status !== "unavailable") || + !(value.reason === null || ( + typeof value.reason === "string" && runtimeObservationReasons.has(value.reason) + )) || + !isNonnegativeInteger(value.resolved_at) + ) return invalidRuntimeObservation(); + + const allocationId = value.instance.allocation_id === null ? null : canonicalUuid(value.instance.allocation_id); + const deviceId = value.instance.device_id === null ? null : canonicalUuid(value.instance.device_id); + const connectionGeneration = value.instance.connection_generation === null + ? null + : canonicalUuid(value.instance.connection_generation); + if ( + (value.instance.allocation_id !== null && allocationId === null) || + (value.instance.device_id !== null && deviceId === null) || + (value.instance.connection_generation !== null && connectionGeneration === null) + ) return invalidRuntimeObservation(); + + const allocationCreatedAt = nullableRuntimeInteger(value.allocation_created_at); + const observedAt = nullableRuntimeInteger(value.observed_at); + const startedAt = nullableRuntimeInteger(value.started_at); + const isNone = value.mode === "none"; + const isSelfHosted = value.mode === "self_hosted"; + const isManaged = value.mode === "openai_hosted"; + if ( + (isNone && ( + value.instance.kind !== "none" || environmentId !== null || value.provider_type !== null || + allocationId !== null || deviceId !== null || connectionGeneration !== null || allocationCreatedAt !== null + )) || + (isSelfHosted && ( + value.instance.kind !== "self_hosted_connection" || environmentId === null || + allocationId !== null || allocationCreatedAt !== null + )) || + (isManaged && ( + value.instance.kind !== "managed_allocation" || environmentId === null || connectionGeneration !== null || + (allocationId === null && (deviceId !== null || allocationCreatedAt !== null)) + )) + ) return invalidRuntimeObservation(); + + const observed = value.status === "observed"; + if ( + (observed && ( + !isManaged || allocationId === null || value.reason !== null || observedAt === null || + observedAt > value.resolved_at + )) || + (!observed && ( + observedAt !== null || startedAt !== null || value.cpu !== null || value.memory !== null + )) || + (value.status === "unsupported" && ( + (!isNone && !isSelfHosted) || value.reason !== "runtime_mode_not_observable" + )) || + (value.status === "unavailable" && ( + !isManaged || value.reason === null || value.reason === "runtime_mode_not_observable" + )) || + (startedAt !== null && observedAt !== null && startedAt > observedAt) || + (allocationCreatedAt !== null && allocationCreatedAt > value.resolved_at) + ) return invalidRuntimeObservation(); + + let cpu: RuntimeObservation["cpu"] = null; + if (value.cpu !== null) { + if (!observed || !isRecord(value.cpu) || !exactFields(value.cpu, runtimeCPUFields)) { + return invalidRuntimeObservation(); + } + cpu = { + usage_seconds_total: nullableRuntimeNumber(value.cpu.usage_seconds_total), + capacity_cores: nullableRuntimeNumber(value.cpu.capacity_cores), + usage_cores: nullableRuntimeNumber(value.cpu.usage_cores), + utilization_ratio: nullableRuntimeNumber(value.cpu.utilization_ratio), + }; + if ( + Object.values(cpu).every((entry) => entry === null) || + (cpu.capacity_cores !== null && cpu.capacity_cores === 0) + ) return invalidRuntimeObservation(); + } + + let memory: RuntimeObservation["memory"] = null; + if (value.memory !== null) { + if (!observed || !isRecord(value.memory) || !exactFields(value.memory, runtimeMemoryFields)) { + return invalidRuntimeObservation(); + } + memory = { + usage_bytes: nullableRuntimeInteger(value.memory.usage_bytes), + limit_bytes: nullableRuntimeInteger(value.memory.limit_bytes), + }; + if ( + (memory.usage_bytes === null && memory.limit_bytes === null) || + memory.limit_bytes === 0 + ) return invalidRuntimeObservation(); + } + + return { + id, object: "agent.runtime_observation", session_id: sessionId, environment_id: environmentId, + mode: value.mode, provider_type: value.provider_type, instance: { + kind: value.instance.kind as RuntimeObservation["instance"]["kind"], + allocation_id: allocationId, device_id: deviceId, connection_generation: connectionGeneration, + }, + status: value.status, reason: value.reason as RuntimeObservation["reason"], + allocation_created_at: allocationCreatedAt, resolved_at: value.resolved_at, + observed_at: observedAt, started_at: startedAt, cpu, memory, + } as RuntimeObservation; +} + +function projectRuntimeObservationList(value: unknown, options?: PageOptions): RuntimeObservationList { + if ( + !isRecord(value) || !exactFields(value, vaultListFields) || value.object !== "list" || + !Array.isArray(value.data) || typeof value.has_more !== "boolean" + ) return invalidRuntimeObservation("Agent Core returned an invalid Runtime observation list."); + const limit = options?.limit ?? 20; + if ( + !Number.isSafeInteger(limit) || limit < 1 || limit > 100 || + (options?.order !== undefined && options.order !== "asc" && options.order !== "desc") || + value.data.length > limit + ) return invalidRuntimeObservation("Agent Core returned an invalid Runtime observation list."); + const data = value.data.map((entry) => projectRuntimeObservation(entry)); + const firstId = data[0]?.id ?? null; + const lastId = data[data.length - 1]?.id ?? null; + if ( + new Set(data.map((entry) => entry.id)).size !== data.length || + value.first_id !== firstId || value.last_id !== lastId || + (value.has_more && data.length === 0) + ) return invalidRuntimeObservation("Agent Core returned an invalid Runtime observation list."); + return { object: "list", data, has_more: value.has_more, first_id: firstId, last_id: lastId }; +} + function projectStreamError(value: unknown): StreamError { if ( !isRecord(value) || !exactFields(value, streamErrorFields) || @@ -2016,6 +2188,31 @@ export class OpenAIAgentsClient implements AgentCore { return { ...page, data: page.data.map((session) => projectAgentSession(session)) }; } + async listRuntimeObservations(options?: PageOptions): Promise { + if ( + (options?.after !== undefined && canonicalUuid(options.after) === null) || + (options?.limit !== undefined && ( + !Number.isSafeInteger(options.limit) || options.limit < 1 || options.limit > 100 + )) || + (options?.order !== undefined && options.order !== "asc" && options.order !== "desc") + ) throw new TypeError("Runtime observation pagination options are invalid."); + const params = new URLSearchParams(); + addPageOptions(params, options); + const value = await this.request( + withQuery("/agents/runtime-observations", params), + { signal: options?.signal }, + ); + return projectRuntimeObservationList(value, options); + } + + async retrieveRuntimeObservation(sessionId: string, options?: ReadOptions): Promise { + const value = await this.request( + `/agents/sessions/${encodeURIComponent(sessionId)}/runtime-observation`, + { signal: options?.signal }, + ); + return projectRuntimeObservation(value, sessionId); + } + async createSession(input: CreateSessionInput, idempotencyKey = createIdempotencyKey()): Promise { if ((input as { stream?: boolean }).stream === true) { throw new TypeError("createSession only supports the JSON response; connect streamEvents after creation."); diff --git a/packages/agents-client/src/protocol-types.test.ts b/packages/agents-client/src/protocol-types.test.ts index 3c1384129..ddf80f36e 100644 --- a/packages/agents-client/src/protocol-types.test.ts +++ b/packages/agents-client/src/protocol-types.test.ts @@ -29,6 +29,8 @@ import type { OpenAIHostedAgentEnvironmentInput, OpenAIHostedAgentEnvironmentResource, RequiredAction, + RuntimeObservation, + RuntimeUnavailableReason, SelfHostedAgentEnvironment, SavedAgentToolInput, SourceFile, @@ -48,6 +50,25 @@ import type { VaultCredential, } from "./types"; +describe("Runtime Observation discriminated contract", () => { + it("narrows status, reason, mode, instance, and sample presence together", () => { + type Observed = Extract; + type Unavailable = Extract; + type NoneMode = Extract; + type SelfHosted = Extract; + + expectTypeOf().toEqualTypeOf<"openai_hosted">(); + expectTypeOf().toEqualTypeOf(); + expectTypeOf().toEqualTypeOf(); + expectTypeOf().toEqualTypeOf(); + expectTypeOf().toEqualTypeOf(); + expectTypeOf().toEqualTypeOf(); + expectTypeOf().toEqualTypeOf(); + expectTypeOf().toEqualTypeOf<"none">(); + expectTypeOf().toEqualTypeOf<"self_hosted_connection">(); + }); +}); + describe("Parsar dadf64a7 basic managed Environment profile", () => { it("pins omitted/default, explicit-enabled, and explicit-disabled network input", () => { const inputs = hostedDadf64.inputs as Record; diff --git a/packages/agents-client/src/types.ts b/packages/agents-client/src/types.ts index 74a156353..8bc3b0f6b 100644 --- a/packages/agents-client/src/types.ts +++ b/packages/agents-client/src/types.ts @@ -682,6 +682,119 @@ export interface CreateSessionStreamOptions extends StreamOptions { onSession: (session: AgentSession) => void; } +export type RuntimeObservationStatus = "observed" | "unsupported" | "unavailable"; +export type RuntimeObservationReason = + | "runtime_mode_not_observable" + | "allocation_pending" + | "runtime_not_running" + | "source_not_configured" + | "sample_timeout" + | "sample_unavailable"; + +export type RuntimeUnavailableReason = Exclude; + +export interface RuntimeCPUObservation { + usage_seconds_total: number | null; + capacity_cores: number | null; + usage_cores: number | null; + utilization_ratio: number | null; +} + +export interface RuntimeMemoryObservation { + usage_bytes: number | null; + limit_bytes: number | null; +} + +interface RuntimeObservationBase { + id: string; + object: "agent.runtime_observation"; + session_id: string; + resolved_at: number; +} + +export interface RuntimeObservedObservation extends RuntimeObservationBase { + environment_id: string; + mode: "openai_hosted"; + provider_type: string | null; + instance: { + kind: "managed_allocation"; + allocation_id: string; + device_id: string | null; + connection_generation: null; + }; + status: "observed"; + reason: null; + allocation_created_at: number | null; + observed_at: number; + started_at: number | null; + cpu: RuntimeCPUObservation | null; + memory: RuntimeMemoryObservation | null; +} + +export interface RuntimeUnavailableObservation extends RuntimeObservationBase { + environment_id: string; + mode: "openai_hosted"; + provider_type: string | null; + instance: { + kind: "managed_allocation"; + allocation_id: string | null; + device_id: string | null; + connection_generation: null; + }; + status: "unavailable"; + reason: RuntimeUnavailableReason; + allocation_created_at: number | null; + observed_at: null; + started_at: null; + cpu: null; + memory: null; +} + +export interface RuntimeNoneObservation extends RuntimeObservationBase { + environment_id: null; + mode: "none"; + provider_type: null; + instance: { kind: "none"; allocation_id: null; device_id: null; connection_generation: null }; + status: "unsupported"; + reason: "runtime_mode_not_observable"; + allocation_created_at: null; + observed_at: null; + started_at: null; + cpu: null; + memory: null; +} + +export interface RuntimeSelfHostedObservation extends RuntimeObservationBase { + environment_id: string; + mode: "self_hosted"; + provider_type: string | null; + instance: { + kind: "self_hosted_connection"; + allocation_id: null; + device_id: string | null; + connection_generation: string | null; + }; + status: "unsupported"; + reason: "runtime_mode_not_observable"; + allocation_created_at: null; + observed_at: null; + started_at: null; + cpu: null; + memory: null; +} + +export type RuntimeObservation = + | RuntimeObservedObservation + | RuntimeUnavailableObservation + | RuntimeNoneObservation + | RuntimeSelfHostedObservation; + +export interface RuntimeObservationList extends ListPage { + object: "list"; + first_id: string | null; + last_id: string | null; +} + export interface AgentCore { listAgents(options?: PageOptions): Promise>; createAgent(input: CreateAgentInput): Promise; @@ -698,6 +811,8 @@ export interface AgentCore { replaceVaultCredentialToken(vaultId: string, credentialId: string, input: ReplaceVaultCredentialTokenInput): Promise; deleteVaultCredential(vaultId: string, credentialId: string): Promise; listSessions(options?: PageOptions & { agentId?: string }): Promise>; + listRuntimeObservations(options?: PageOptions): Promise; + retrieveRuntimeObservation(sessionId: string, options?: ReadOptions): Promise; createSession(input: CreateSessionInput, idempotencyKey?: string): Promise; createSessionStream( input: Omit, diff --git a/services/agents-api/cmd/server/main.go b/services/agents-api/cmd/server/main.go index 757c2be80..d7c2daf88 100644 --- a/services/agents-api/cmd/server/main.go +++ b/services/agents-api/cmd/server/main.go @@ -28,6 +28,7 @@ import ( "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/execution" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtime" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimeenrollment" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimeobs" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/store" "github.com/jackc/pgx/v5/pgxpool" ) @@ -89,9 +90,27 @@ func run() error { if err := executionStore.EnsureProjectScopes(ready, auth.ProjectScopes()); err != nil { return err } + observationSources := map[string]runtimeobs.Source{} + if managed != nil { + for key, provider := range managed.Providers { + source, ok := provider.(runtimeobs.Source) + if !ok { + continue + } + observationSources[key] = source + } + } + resolver, err := runtimeobs.NewResolver(executionStore) + if err != nil { + return err + } + observationService, err := runtimeobs.NewService(resolver, observationSources) + if err != nil { + return err + } var workerDone chan error var worker *execution.Worker - options := []api.Option{api.WithSubagents(executionStore), api.WithSkills(executionStore), api.WithSourceFiles(executionStore), api.WithSessionArtifacts(executionStore)} + options := []api.Option{api.WithSubagents(executionStore), api.WithSkills(executionStore), api.WithSourceFiles(executionStore), api.WithSessionArtifacts(executionStore), api.WithRuntimeObservations(observationService)} var daemonHandler http.Handler var registry *gateway.Registry if wsURL := os.Getenv("AGENTS_API_DAEMON_WS_URL"); wsURL != "" { diff --git a/services/agents-api/internal/api/handler.go b/services/agents-api/internal/api/handler.go index 4fa78bdbb..a154843ef 100644 --- a/services/agents-api/internal/api/handler.go +++ b/services/agents-api/internal/api/handler.go @@ -35,20 +35,21 @@ type ResourceStore interface { } type Handler struct { - policy execution.Policy - store ResourceStore - auth *Authenticator - harnesses map[string]bool - engine string - inputs InputSubmitter - executorURL string - hostedEnvironments bool - directoryReader EnvironmentDirectoryReader - fileWriter EnvironmentFileWriter - skills SkillStore - sourceFiles SourceFileStore - artifacts SessionArtifactStore - subagents SubagentStore + policy execution.Policy + store ResourceStore + auth *Authenticator + harnesses map[string]bool + engine string + inputs InputSubmitter + executorURL string + hostedEnvironments bool + directoryReader EnvironmentDirectoryReader + fileWriter EnvironmentFileWriter + skills SkillStore + sourceFiles SourceFileStore + artifacts SessionArtifactStore + subagents SubagentStore + runtimeObservations RuntimeObservationService } func NewHandler(s ResourceStore, auth *Authenticator, engine string, options ...Option) (http.Handler, error) { @@ -100,6 +101,8 @@ func NewHandler(s ResourceStore, auth *Authenticator, engine string, options ... r.Post("/agents/sessions", h.createSession) r.Get("/agents/sessions", h.listSessions) r.Get("/agents/sessions/{session_id}", h.getSession) + r.Get("/agents/sessions/{session_id}/runtime-observation", h.getRuntimeObservation) + r.Get("/agents/runtime-observations", h.listRuntimeObservations) r.Post("/agents/sessions/{session_id}", h.updateSession) r.Delete("/agents/sessions/{session_id}", h.deleteSession) r.Post("/agents/sessions/{session_id}/events", h.createEvents) diff --git a/services/agents-api/internal/api/handler_test.go b/services/agents-api/internal/api/handler_test.go index dbf08bcf7..b1c3befe2 100644 --- a/services/agents-api/internal/api/handler_test.go +++ b/services/agents-api/internal/api/handler_test.go @@ -19,8 +19,19 @@ import ( type recordingStore struct { ResourceStore - tenant string - input store.CreateSessionInput + tenant string + input store.CreateSessionInput + sessions []store.Session + nextSessionCursor string + listTenant string + listAfter string + listLimit int + listAscending bool +} + +func (s *recordingStore) ListSessions(_ context.Context, tenant, after string, limit int, ascending bool, _ *string) (store.SessionPage, error) { + s.listTenant, s.listAfter, s.listLimit, s.listAscending = tenant, after, limit, ascending + return store.SessionPage{Sessions: append([]store.Session(nil), s.sessions...), NextCursor: s.nextSessionCursor}, nil } func (s *recordingStore) GetSession(ctx context.Context, tenant, id string) (store.Session, error) { diff --git a/services/agents-api/internal/api/runtime_observations.go b/services/agents-api/internal/api/runtime_observations.go new file mode 100644 index 000000000..54d061a2d --- /dev/null +++ b/services/agents-api/internal/api/runtime_observations.go @@ -0,0 +1,221 @@ +package api + +import ( + "context" + "errors" + "net/http" + "regexp" + "sync" + "time" + + v1 "github.com/MiniMax-AI-Dev/parsar/contracts/agents-api/v1" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimeobs" + "github.com/go-chi/chi/v5" +) + +const ( + runtimeObservationConcurrency = 8 + runtimeObservationSourceBudget = 2 * time.Second + runtimeObservationRequestBudget = 10 * time.Second +) + +var runtimeProviderTypePattern = regexp.MustCompile(`^[a-z][a-z0-9_]{0,31}$`) + +type RuntimeObservationService interface { + ObserveSession(context.Context, string, string) (runtimeobs.Observation, error) +} + +func WithRuntimeObservations(service RuntimeObservationService) Option { + return func(h *Handler) { h.runtimeObservations = service } +} + +// @Summary Retrieve a Session Runtime observation +// @Description Core extension returning one tenant-scoped, read-only current Runtime observation. It never provisions, renews, restarts, pauses or stops compute. +// @Tags Runtime observations +// @Produce json +// @Security BearerAuth +// @Param OpenAI-Beta header string true "agents=v1" +// @Param session_id path string true "Session ID" +// @Success 200 {object} v1.RuntimeObservation +// @Failure 400,401,404,500,503 {object} v1.ErrorResponse +// @Router /agents/sessions/{session_id}/runtime-observation [get] +func (h *Handler) getRuntimeObservation(w http.ResponseWriter, r *http.Request) { + if len(r.URL.Query()) != 0 { + writeError(w, http.StatusBadRequest, "unsupported_parameter", "Runtime observation retrieval does not accept query parameters.") + return + } + if h.runtimeObservations == nil { + writeError(w, http.StatusServiceUnavailable, "execution_unavailable", "Runtime observation is not configured on this service.") + return + } + ctx, cancel := context.WithTimeout(r.Context(), runtimeObservationSourceBudget) + defer cancel() + observation, err := h.runtimeObservations.ObserveSession(ctx, tenantID(r), chi.URLParam(r, "session_id")) + if err != nil { + writeStoreError(w, r, err) + return + } + response, err := runtimeObservationResponse(observation) + if err != nil { + writeStoreError(w, r, err) + return + } + writeJSON(w, http.StatusOK, response) +} + +// @Summary List current Runtime observations +// @Description Core extension listing one current Runtime context per tenant-owned Session in Session creation order. Each row has an independent resolved_at and optional provider observed_at; the page is not an atomic telemetry snapshot. +// @Tags Runtime observations +// @Produce json +// @Security BearerAuth +// @Param OpenAI-Beta header string true "agents=v1" +// @Param after query string false "Last observation ID from the previous page" +// @Param limit query int false "Page size" minimum(1) maximum(100) default(20) +// @Param order query string false "Session creation order" Enums(asc,desc) default(desc) +// @Success 200 {object} v1.RuntimeObservationList +// @Failure 400,401,404,500,503 {object} v1.ErrorResponse +// @Router /agents/runtime-observations [get] +func (h *Handler) listRuntimeObservations(w http.ResponseWriter, r *http.Request) { + if h.runtimeObservations == nil { + writeError(w, http.StatusServiceUnavailable, "execution_unavailable", "Runtime observation is not configured on this service.") + return + } + options, ok := readPage(w, r) + if !ok { + return + } + ctx, cancel := context.WithTimeout(r.Context(), runtimeObservationRequestBudget) + defer cancel() + page, err := h.store.ListSessions(ctx, tenantID(r), options.after, options.limit, options.ascending, nil) + if err != nil { + writeStoreError(w, r, err) + return + } + observations := make([]runtimeobs.Observation, len(page.Sessions)) + semaphore := make(chan struct{}, runtimeObservationConcurrency) + work, stop := context.WithCancel(ctx) + defer stop() + var wait sync.WaitGroup + var once sync.Once + var firstErr error + for index, session := range page.Sessions { + wait.Add(1) + go func(index int, sessionID string) { + defer wait.Done() + select { + case semaphore <- struct{}{}: + defer func() { <-semaphore }() + case <-work.Done(): + return + } + sampleCtx, sampleCancel := context.WithTimeout(work, runtimeObservationSourceBudget) + defer sampleCancel() + value, err := h.runtimeObservations.ObserveSession(sampleCtx, tenantID(r), sessionID) + if err != nil { + once.Do(func() { firstErr = err; stop() }) + return + } + observations[index] = value + }(index, session.ID) + } + wait.Wait() + if firstErr != nil { + writeStoreError(w, r, firstErr) + return + } + if err := ctx.Err(); err != nil { + writeError(w, http.StatusServiceUnavailable, "execution_unavailable", "Runtime observation collection exceeded its request budget.") + return + } + response := v1.RuntimeObservationList{Object: "list", Data: make([]v1.RuntimeObservation, 0, len(observations)), HasMore: page.NextCursor != ""} + for _, observation := range observations { + item, err := runtimeObservationResponse(observation) + if err != nil { + writeStoreError(w, r, err) + return + } + response.Data = append(response.Data, item) + } + if len(response.Data) > 0 { + response.FirstID = &response.Data[0].ID + response.LastID = &response.Data[len(response.Data)-1].ID + } + writeJSON(w, http.StatusOK, response) +} + +func runtimeObservationResponse(observation runtimeobs.Observation) (v1.RuntimeObservation, error) { + if observation.Target.SessionID == "" || observation.ResolvedAt.IsZero() || observation.ResolvedAt.Unix() < 0 { + return v1.RuntimeObservation{}, errors.New("invalid Runtime observation identity") + } + if !observation.Target.Instance.AllocationCreatedAt.IsZero() && + (observation.Target.Instance.AllocationCreatedAt.Unix() < 0 || observation.Target.Instance.AllocationCreatedAt.After(observation.ResolvedAt)) { + return v1.RuntimeObservation{}, errors.New("invalid Runtime allocation creation time") + } + if observation.Sample != nil { + if observation.Sample.ObservedAt.IsZero() || observation.Sample.ObservedAt.Unix() < 0 || observation.Sample.ObservedAt.After(observation.ResolvedAt) { + return v1.RuntimeObservation{}, errors.New("invalid Runtime sample time") + } + if observation.Sample.StartedAt != nil && + (observation.Sample.StartedAt.IsZero() || observation.Sample.StartedAt.Unix() < 0 || observation.Sample.StartedAt.After(observation.Sample.ObservedAt)) { + return v1.RuntimeObservation{}, errors.New("invalid Runtime start time") + } + } + result := v1.RuntimeObservation{ + ID: observation.Target.SessionID, Object: "agent.runtime_observation", SessionID: observation.Target.SessionID, + Mode: string(observation.Target.Mode), Status: string(observation.Status), ResolvedAt: observation.ResolvedAt.Unix(), + } + if observation.Target.EnvironmentID != "" { + result.EnvironmentID = &observation.Target.EnvironmentID + } + if observation.ProviderType != "" { + if !runtimeProviderTypePattern.MatchString(observation.ProviderType) { + return v1.RuntimeObservation{}, errors.New("invalid Runtime observation provider type") + } + result.ProviderType = &observation.ProviderType + } + if observation.Reason != "" { + result.Reason = &observation.Reason + } + switch observation.Target.Mode { + case runtimeobs.ModeManaged: + result.Instance.Kind = "managed_allocation" + if observation.Target.Instance.AllocationID != "" { + result.Instance.AllocationID = &observation.Target.Instance.AllocationID + } + if observation.Target.Instance.DeviceID != "" { + result.Instance.DeviceID = &observation.Target.Instance.DeviceID + } + if !observation.Target.Instance.AllocationCreatedAt.IsZero() { + created := observation.Target.Instance.AllocationCreatedAt.Unix() + result.AllocationCreatedAt = &created + } + case runtimeobs.ModeSelfHosted: + result.Instance.Kind = "self_hosted_connection" + if observation.Target.Instance.DeviceID != "" { + result.Instance.DeviceID = &observation.Target.Instance.DeviceID + } + if observation.Target.Instance.ConnectionGeneration != "" { + result.Instance.ConnectionGeneration = &observation.Target.Instance.ConnectionGeneration + } + case runtimeobs.ModeNone: + result.Instance.Kind = "none" + default: + return v1.RuntimeObservation{}, errors.New("invalid Runtime observation mode") + } + if observation.Sample == nil { + return result, nil + } + observedAt := observation.Sample.ObservedAt.Unix() + result.ObservedAt = &observedAt + if observation.Sample.StartedAt != nil { + startedAt := observation.Sample.StartedAt.Unix() + result.StartedAt = &startedAt + } + if observation.Sample.CPUUsageSecondsTotal != nil || observation.Sample.CPUCapacityCores != nil { + result.CPU = &v1.RuntimeCPUObservation{UsageSecondsTotal: observation.Sample.CPUUsageSecondsTotal, CapacityCores: observation.Sample.CPUCapacityCores} + } + if observation.Sample.MemoryUsageBytes != nil || observation.Sample.MemoryLimitBytes != nil { + result.Memory = &v1.RuntimeMemoryObservation{UsageBytes: observation.Sample.MemoryUsageBytes, LimitBytes: observation.Sample.MemoryLimitBytes} + } + return result, nil +} diff --git a/services/agents-api/internal/api/runtime_observations_test.go b/services/agents-api/internal/api/runtime_observations_test.go new file mode 100644 index 000000000..29768be94 --- /dev/null +++ b/services/agents-api/internal/api/runtime_observations_test.go @@ -0,0 +1,224 @@ +package api + +import ( + "context" + "encoding/json" + "errors" + "net/http" + "net/http/httptest" + "sync/atomic" + "testing" + "time" + + v1 "github.com/MiniMax-AI-Dev/parsar/contracts/agents-api/v1" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimeobs" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/store" + "github.com/google/uuid" +) + +type runtimeObservationFixture struct { + values map[string]runtimeobs.Observation +} + +func (f runtimeObservationFixture) ObserveSession(_ context.Context, tenant, session string) (runtimeobs.Observation, error) { + value := f.values[session] + value.Target.TenantID = tenant + return value, nil +} + +type runtimeObservationServiceFunc func(context.Context, string, string) (runtimeobs.Observation, error) + +func (f runtimeObservationServiceFunc) ObserveSession(ctx context.Context, tenant, session string) (runtimeobs.Observation, error) { + return f(ctx, tenant, session) +} + +func runtimeObservationRequest(handler http.Handler, path string) *httptest.ResponseRecorder { + request := httptest.NewRequest(http.MethodGet, path, nil) + request.Header.Set("Authorization", "Bearer test-api-key") + request.Header.Set("OpenAI-Beta", "agents=v1") + response := httptest.NewRecorder() + handler.ServeHTTP(response, request) + return response +} + +func TestRuntimeObservationRoutesUseSessionIdentityAndExactNullability(t *testing.T) { + now := time.Date(2026, 9, 22, 8, 0, 0, 0, time.UTC) + sessionID := uuid.NewString() + service := runtimeObservationFixture{values: map[string]runtimeobs.Observation{ + sessionID: {Target: runtimeobs.Target{SessionID: sessionID, Mode: runtimeobs.ModeNone}, Status: runtimeobs.StatusUnsupported, Reason: "runtime_mode_not_observable", ResolvedAt: now}, + }} + handler, saved, _ := testHandler(t, WithRuntimeObservations(service)) + saved.sessions = []store.Session{{ID: sessionID, CreatedAt: now}} + + for _, path := range []string{"/v1/agents/sessions/" + sessionID + "/runtime-observation", "/v1/agents/runtime-observations?limit=1"} { + response := runtimeObservationRequest(handler, path) + if response.Code != http.StatusOK { + t.Fatalf("%s returned %d: %s", path, response.Code, response.Body) + } + if path[len(path)-7:] == "limit=1" { + var page v1.RuntimeObservationList + if json.Unmarshal(response.Body.Bytes(), &page) != nil || len(page.Data) != 1 || page.FirstID == nil || *page.FirstID != sessionID || page.LastID == nil || *page.LastID != sessionID { + t.Fatalf("invalid observation page: %s", response.Body) + } + continue + } + var value v1.RuntimeObservation + if json.Unmarshal(response.Body.Bytes(), &value) != nil || value.ID != sessionID || value.SessionID != sessionID || value.Instance.Kind != "none" || value.EnvironmentID != nil || value.ProviderType != nil || value.ObservedAt != nil || value.CPU != nil || value.Memory != nil || value.Reason == nil || *value.Reason != "runtime_mode_not_observable" { + t.Fatalf("invalid unsupported observation: %s", response.Body) + } + } +} + +func TestRuntimeObservationRoutesRequireConfiguredServiceAndRejectQueries(t *testing.T) { + handler, _, _ := testHandler(t) + missing := runtimeObservationRequest(handler, "/v1/agents/runtime-observations") + if missing.Code != http.StatusServiceUnavailable { + t.Fatalf("unconfigured service returned %d: %s", missing.Code, missing.Body) + } + + sessionID := uuid.NewString() + now := time.Now().UTC() + service := runtimeObservationFixture{values: map[string]runtimeobs.Observation{ + sessionID: {Target: runtimeobs.Target{SessionID: sessionID, Mode: runtimeobs.ModeNone}, Status: runtimeobs.StatusUnsupported, Reason: "runtime_mode_not_observable", ResolvedAt: now}, + }} + handler, _, _ = testHandler(t, WithRuntimeObservations(service)) + invalid := runtimeObservationRequest(handler, "/v1/agents/sessions/"+sessionID+"/runtime-observation?provider=docker") + if invalid.Code != http.StatusBadRequest { + t.Fatalf("unsupported query returned %d: %s", invalid.Code, invalid.Body) + } +} + +func TestRuntimeObservationResponsePreservesObservedZero(t *testing.T) { + zeroCPU := float64(0) + zeroMemory := uint64(0) + now := time.Date(2026, 9, 22, 8, 0, 0, 0, time.UTC) + sessionID, environmentID := uuid.NewString(), uuid.NewString() + value, err := runtimeObservationResponse(runtimeobs.Observation{ + Target: runtimeobs.Target{SessionID: sessionID, EnvironmentID: environmentID, Mode: runtimeobs.ModeManaged, Instance: runtimeobs.Instance{AllocationID: uuid.NewString(), DeviceID: uuid.NewString(), AllocationCreatedAt: now.Add(-time.Hour)}}, + Status: runtimeobs.StatusObserved, ProviderType: "docker", ResolvedAt: now, + Sample: &runtimeobs.Sample{ObservedAt: now, CPUUsageSecondsTotal: &zeroCPU, MemoryUsageBytes: &zeroMemory}, + }) + if err != nil || value.CPU == nil || value.CPU.UsageSecondsTotal == nil || *value.CPU.UsageSecondsTotal != 0 || value.Memory == nil || value.Memory.UsageBytes == nil || *value.Memory.UsageBytes != 0 { + t.Fatalf("observed zero was lost: %+v %v", value, err) + } +} + +func TestRuntimeObservationResponseRejectsTimesOutsidePublicContract(t *testing.T) { + now := time.Date(2026, 9, 22, 8, 0, 0, 0, time.UTC) + preEpoch := time.Unix(-1, 0).UTC() + base := runtimeobs.Observation{ + Target: runtimeobs.Target{ + SessionID: uuid.NewString(), EnvironmentID: uuid.NewString(), Mode: runtimeobs.ModeManaged, + Instance: runtimeobs.Instance{AllocationID: uuid.NewString(), DeviceID: uuid.NewString(), AllocationCreatedAt: now.Add(-time.Hour)}, + }, + Status: runtimeobs.StatusObserved, ResolvedAt: now, + Sample: &runtimeobs.Sample{ObservedAt: now, StartedAt: timePointer(now.Add(-time.Minute))}, + } + for _, mutate := range []func(*runtimeobs.Observation){ + func(value *runtimeobs.Observation) { value.Target.Instance.AllocationCreatedAt = now.Add(time.Second) }, + func(value *runtimeobs.Observation) { value.Target.Instance.AllocationCreatedAt = preEpoch }, + func(value *runtimeobs.Observation) { value.Sample.ObservedAt = preEpoch }, + func(value *runtimeobs.Observation) { value.Sample.StartedAt = &preEpoch }, + } { + observation := base + sample := *base.Sample + observation.Sample = &sample + mutate(&observation) + if _, err := runtimeObservationResponse(observation); err == nil { + t.Fatalf("invalid Runtime time accepted: %+v", observation) + } + } +} + +func timePointer(value time.Time) *time.Time { return &value } + +func TestRuntimeObservationListPreservesStoreOrderAndTenantPagination(t *testing.T) { + now := time.Date(2026, 9, 22, 8, 0, 0, 0, time.UTC) + sessionIDs := []string{uuid.NewString(), uuid.NewString(), uuid.NewString()} + var expectedTenant string + service := runtimeObservationServiceFunc(func(_ context.Context, tenant, session string) (runtimeobs.Observation, error) { + if tenant != expectedTenant { + return runtimeobs.Observation{}, errors.New("unexpected tenant") + } + if session == sessionIDs[0] { + time.Sleep(20 * time.Millisecond) + } + return runtimeobs.Observation{ + Target: runtimeobs.Target{SessionID: session, Mode: runtimeobs.ModeNone}, + Status: runtimeobs.StatusUnsupported, Reason: "runtime_mode_not_observable", ResolvedAt: now, + }, nil + }) + handler, saved, tenant := testHandler(t, WithRuntimeObservations(service)) + expectedTenant = tenant + for _, id := range sessionIDs { + saved.sessions = append(saved.sessions, store.Session{ID: id, CreatedAt: now}) + } + saved.nextSessionCursor = "next" + + response := runtimeObservationRequest(handler, "/v1/agents/runtime-observations?after=cursor&limit=3&order=asc") + if response.Code != http.StatusOK { + t.Fatalf("list returned %d: %s", response.Code, response.Body) + } + var page v1.RuntimeObservationList + if err := json.Unmarshal(response.Body.Bytes(), &page); err != nil { + t.Fatal(err) + } + if len(page.Data) != len(sessionIDs) || !page.HasMore || saved.listTenant != tenant || saved.listAfter != "cursor" || saved.listLimit != 3 || !saved.listAscending { + t.Fatalf("pagination binding was not preserved: page=%+v store=%+v", page, saved) + } + for index, item := range page.Data { + if item.ID != sessionIDs[index] { + t.Fatalf("concurrent collection reordered page: %+v", page.Data) + } + } +} + +func TestRuntimeObservationListBoundsCollectionConcurrency(t *testing.T) { + now := time.Date(2026, 9, 22, 8, 0, 0, 0, time.UTC) + var active, maximum atomic.Int32 + service := runtimeObservationServiceFunc(func(_ context.Context, _, session string) (runtimeobs.Observation, error) { + current := active.Add(1) + defer active.Add(-1) + for current > maximum.Load() && !maximum.CompareAndSwap(maximum.Load(), current) { + } + time.Sleep(15 * time.Millisecond) + return runtimeobs.Observation{ + Target: runtimeobs.Target{SessionID: session, Mode: runtimeobs.ModeNone}, + Status: runtimeobs.StatusUnsupported, Reason: "runtime_mode_not_observable", ResolvedAt: now, + }, nil + }) + handler, saved, _ := testHandler(t, WithRuntimeObservations(service)) + for range 20 { + saved.sessions = append(saved.sessions, store.Session{ID: uuid.NewString(), CreatedAt: now}) + } + response := runtimeObservationRequest(handler, "/v1/agents/runtime-observations?limit=20") + if response.Code != http.StatusOK { + t.Fatalf("list returned %d: %s", response.Code, response.Body) + } + if got := maximum.Load(); got == 0 || got > runtimeObservationConcurrency { + t.Fatalf("collection concurrency = %d, want 1..%d", got, runtimeObservationConcurrency) + } +} + +func TestRuntimeObservationListRejectsWholePageOnIntegrityFailure(t *testing.T) { + now := time.Date(2026, 9, 22, 8, 0, 0, 0, time.UTC) + validID, invalidID := uuid.NewString(), uuid.NewString() + service := runtimeObservationFixture{values: map[string]runtimeobs.Observation{ + validID: { + Target: runtimeobs.Target{SessionID: validID, Mode: runtimeobs.ModeNone}, + Status: runtimeobs.StatusUnsupported, Reason: "runtime_mode_not_observable", ResolvedAt: now, + }, + invalidID: {Target: runtimeobs.Target{Mode: runtimeobs.ModeNone}, Status: runtimeobs.StatusUnsupported, ResolvedAt: now}, + }} + handler, saved, _ := testHandler(t, WithRuntimeObservations(service)) + saved.sessions = []store.Session{{ID: validID, CreatedAt: now}, {ID: invalidID, CreatedAt: now}} + + response := runtimeObservationRequest(handler, "/v1/agents/runtime-observations?limit=2") + if response.Code != http.StatusInternalServerError { + t.Fatalf("integrity failure returned %d: %s", response.Code, response.Body) + } + var envelope v1.ErrorResponse + if err := json.Unmarshal(response.Body.Bytes(), &envelope); err != nil || envelope.Error.Code == "" { + t.Fatalf("integrity failure leaked a partial page: %s", response.Body) + } +} diff --git a/services/agents-api/internal/runtimeobs/identity.go b/services/agents-api/internal/runtimeobs/identity.go index 7e33f4bb9..17e38f7d4 100644 --- a/services/agents-api/internal/runtimeobs/identity.go +++ b/services/agents-api/internal/runtimeobs/identity.go @@ -1,5 +1,7 @@ package runtimeobs +import "time" + // Instance is one provider-owned Runtime incarnation. AllocationID is present // for managed compute. DeviceID and ConnectionGeneration are reserved for a // future authenticated self-hosted telemetry source. @@ -8,6 +10,8 @@ type Instance struct { ProviderKey string DeviceID string ConnectionGeneration string + AllocationState string + AllocationCreatedAt time.Time } // Target binds telemetry to durable Core identity. A Session is not itself a diff --git a/services/agents-api/internal/runtimeobs/resolver.go b/services/agents-api/internal/runtimeobs/resolver.go index 7b1047c0c..58a5f28df 100644 --- a/services/agents-api/internal/runtimeobs/resolver.go +++ b/services/agents-api/internal/runtimeobs/resolver.go @@ -9,7 +9,10 @@ import ( "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/store" ) -var ErrUnavailable = errors.New("Runtime observation unavailable") +var ( + ErrUnavailable = errors.New("Runtime observation unavailable") + ErrNotRunning = errors.New("Runtime is not running") +) type sessionStore interface { GetSession(context.Context, string, string) (store.Session, error) @@ -49,12 +52,18 @@ func (r *Resolver) Resolve(ctx context.Context, tenantID, sessionID string) (Tar if session.Environment == nil { return Target{}, errors.New("self-hosted Session is missing its Environment") } + if session.Environment.TenantID != session.TenantID || session.Environment.SessionID != session.ID { + return Target{}, errors.New("self-hosted Environment does not match resolved ownership") + } target.EnvironmentID = session.Environment.ID return target, nil case ModeManaged: if session.Environment == nil { return Target{}, errors.New("managed Session is missing its Environment") } + if session.Environment.TenantID != session.TenantID || session.Environment.SessionID != session.ID { + return Target{}, errors.New("managed Environment does not match resolved ownership") + } target.EnvironmentID = session.Environment.ID allocation, err := r.store.GetRuntimeAllocation(ctx, tenantID, target.EnvironmentID) if errors.Is(err, store.ErrNotFound) { @@ -63,7 +72,13 @@ func (r *Resolver) Resolve(ctx context.Context, tenantID, sessionID string) (Tar if err != nil { return Target{}, fmt.Errorf("resolve Runtime allocation: %w", err) } - target.Instance = Instance{AllocationID: allocation.ID, ProviderKey: allocation.ProviderKey, DeviceID: allocation.DeviceID} + if allocation.TenantID != tenantID || allocation.SessionID != session.ID || allocation.EnvironmentID != target.EnvironmentID { + return Target{}, errors.New("Runtime allocation does not match resolved ownership") + } + target.Instance = Instance{ + AllocationID: allocation.ID, ProviderKey: allocation.ProviderKey, DeviceID: allocation.DeviceID, + AllocationState: allocation.State, AllocationCreatedAt: allocation.CreatedAt, + } return target, nil default: return Target{}, errors.New("invalid stored Runtime environment type") diff --git a/services/agents-api/internal/runtimeobs/resolver_test.go b/services/agents-api/internal/runtimeobs/resolver_test.go index 738c79f2c..7edb4ba48 100644 --- a/services/agents-api/internal/runtimeobs/resolver_test.go +++ b/services/agents-api/internal/runtimeobs/resolver_test.go @@ -24,8 +24,11 @@ func (s resolverStore) GetRuntimeAllocation(context.Context, string, string) (st func TestResolverBindsManagedSessionEnvironmentAndAllocation(t *testing.T) { r, err := NewResolver(resolverStore{ - session: store.Session{ID: "session", TenantID: "tenant", Configuration: []byte(`{"environment":{"type":"openai_hosted"}}`), Environment: &store.Environment{ID: "environment"}}, - allocation: store.RuntimeAllocation{ID: "allocation", ProviderKey: "provider", DeviceID: "device"}, + session: store.Session{ID: "session", TenantID: "tenant", Configuration: []byte(`{"environment":{"type":"openai_hosted"}}`), Environment: &store.Environment{ID: "environment", TenantID: "tenant", SessionID: "session"}}, + allocation: store.RuntimeAllocation{ + ID: "allocation", TenantID: "tenant", SessionID: "session", EnvironmentID: "environment", + ProviderKey: "provider", DeviceID: "device", + }, }) if err != nil { t.Fatal(err) @@ -45,7 +48,7 @@ func TestResolverKeepsUnsupportedModesDistinct(t *testing.T) { environment *store.Environment }{ {mode: "none"}, - {mode: "self_hosted", environment: &store.Environment{ID: "environment"}}, + {mode: "self_hosted", environment: &store.Environment{ID: "environment", TenantID: "tenant", SessionID: "session"}}, } { r, err := NewResolver(resolverStore{session: store.Session{ID: "session", TenantID: "tenant", Configuration: []byte(`{"environment":{"type":"` + tc.mode + `"}}`), Environment: tc.environment}}) if err != nil { @@ -60,7 +63,7 @@ func TestResolverKeepsUnsupportedModesDistinct(t *testing.T) { func TestResolverReportsManagedAllocationAsUnavailable(t *testing.T) { r, err := NewResolver(resolverStore{ - session: store.Session{ID: "session", TenantID: "tenant", Configuration: []byte(`{"environment":{"type":"openai_hosted"}}`), Environment: &store.Environment{ID: "environment"}}, + session: store.Session{ID: "session", TenantID: "tenant", Configuration: []byte(`{"environment":{"type":"openai_hosted"}}`), Environment: &store.Environment{ID: "environment", TenantID: "tenant", SessionID: "session"}}, allocationErr: store.ErrNotFound, }) if err != nil { @@ -71,3 +74,51 @@ func TestResolverReportsManagedAllocationAsUnavailable(t *testing.T) { t.Fatalf("allocation absence was not preserved: %+v %v", target, err) } } + +func TestResolverRejectsMismatchedEnvironmentOwnership(t *testing.T) { + for _, mode := range []string{"self_hosted", "openai_hosted"} { + for _, environment := range []store.Environment{ + {ID: "environment", TenantID: "other", SessionID: "session"}, + {ID: "environment", TenantID: "tenant", SessionID: "other"}, + } { + resolver, err := NewResolver(resolverStore{session: store.Session{ + ID: "session", TenantID: "tenant", + Configuration: []byte(`{"environment":{"type":"` + mode + `"}}`), Environment: &environment, + }}) + if err != nil { + t.Fatal(err) + } + if _, err := resolver.Resolve(t.Context(), "tenant", "session"); err == nil { + t.Fatalf("mismatched %s Environment accepted: %+v", mode, environment) + } + } + } +} + +func TestResolverRejectsMismatchedAllocationOwnership(t *testing.T) { + base := store.RuntimeAllocation{ + ID: "allocation", TenantID: "tenant", SessionID: "session", EnvironmentID: "environment", + ProviderKey: "provider", DeviceID: "device", + } + for _, mutate := range []func(*store.RuntimeAllocation){ + func(value *store.RuntimeAllocation) { value.TenantID = "other" }, + func(value *store.RuntimeAllocation) { value.SessionID = "other" }, + func(value *store.RuntimeAllocation) { value.EnvironmentID = "other" }, + } { + allocation := base + mutate(&allocation) + resolver, err := NewResolver(resolverStore{ + session: store.Session{ + ID: "session", TenantID: "tenant", Configuration: []byte(`{"environment":{"type":"openai_hosted"}}`), + Environment: &store.Environment{ID: "environment", TenantID: "tenant", SessionID: "session"}, + }, + allocation: allocation, + }) + if err != nil { + t.Fatal(err) + } + if _, err := resolver.Resolve(t.Context(), "tenant", "session"); err == nil { + t.Fatalf("mismatched allocation accepted: %+v", allocation) + } + } +} diff --git a/services/agents-api/internal/runtimeobs/sample.go b/services/agents-api/internal/runtimeobs/sample.go index bf653c174..5b36e81f1 100644 --- a/services/agents-api/internal/runtimeobs/sample.go +++ b/services/agents-api/internal/runtimeobs/sample.go @@ -4,6 +4,7 @@ package runtimeobs import ( "errors" + "math" "time" ) @@ -36,19 +37,23 @@ type Sample struct { } func (s Sample) validate(now time.Time) error { - if s.ObservedAt.IsZero() || s.ObservedAt.After(now) { + if s.ObservedAt.IsZero() || s.ObservedAt.Unix() < 0 || s.ObservedAt.After(now) { return errors.New("invalid Runtime observation time") } - if s.StartedAt != nil && (s.StartedAt.IsZero() || s.StartedAt.After(s.ObservedAt)) { + if s.StartedAt != nil && (s.StartedAt.IsZero() || s.StartedAt.Unix() < 0 || s.StartedAt.After(s.ObservedAt)) { return errors.New("invalid Runtime start time") } - if s.CPUUsageSecondsTotal != nil && *s.CPUUsageSecondsTotal < 0 { + if s.CPUUsageSecondsTotal != nil && (*s.CPUUsageSecondsTotal < 0 || math.IsNaN(*s.CPUUsageSecondsTotal) || math.IsInf(*s.CPUUsageSecondsTotal, 0)) { return errors.New("invalid Runtime CPU usage") } - if s.CPUCapacityCores != nil && *s.CPUCapacityCores <= 0 { + if s.CPUCapacityCores != nil && (*s.CPUCapacityCores <= 0 || math.IsNaN(*s.CPUCapacityCores) || math.IsInf(*s.CPUCapacityCores, 0)) { return errors.New("invalid Runtime CPU capacity") } - if s.MemoryLimitBytes != nil && *s.MemoryLimitBytes == 0 { + const maxSafeJSONInteger = uint64(1<<53 - 1) + if s.MemoryUsageBytes != nil && *s.MemoryUsageBytes > maxSafeJSONInteger { + return errors.New("Runtime memory usage exceeds the public JSON integer range") + } + if s.MemoryLimitBytes != nil && (*s.MemoryLimitBytes == 0 || *s.MemoryLimitBytes > maxSafeJSONInteger) { return errors.New("invalid Runtime memory limit") } return nil diff --git a/services/agents-api/internal/runtimeobs/service.go b/services/agents-api/internal/runtimeobs/service.go index e05f7c678..e56d17348 100644 --- a/services/agents-api/internal/runtimeobs/service.go +++ b/services/agents-api/internal/runtimeobs/service.go @@ -8,10 +8,12 @@ import ( ) type Observation struct { - Target Target - Status Status - Sample *Sample - Reason string + Target Target + Status Status + Sample *Sample + Reason string + ProviderType string + ResolvedAt time.Time } type Service struct { @@ -36,25 +38,59 @@ func NewService(resolver TargetResolver, sources map[string]Source) (*Service, e func (s *Service) ObserveSession(ctx context.Context, tenantID, sessionID string) (Observation, error) { target, err := s.resolver.Resolve(ctx, tenantID, sessionID) + resolvedAt := s.now() if errors.Is(err, ErrUnavailable) { - return Observation{Target: target, Status: StatusUnavailable, Reason: "runtime_allocation_unavailable"}, nil + if target.TenantID != tenantID || target.SessionID != sessionID || target.Mode != ModeManaged || target.EnvironmentID == "" { + return Observation{}, errors.New("Runtime observation resolver returned invalid pending allocation identity") + } + return Observation{Target: target, Status: StatusUnavailable, Reason: "allocation_pending", ResolvedAt: resolvedAt}, nil } if err != nil { return Observation{}, err } + if target.TenantID != tenantID || target.SessionID != sessionID { + return Observation{}, errors.New("Runtime observation resolver returned mismatched ownership") + } + if (target.Mode == ModeNone && target.EnvironmentID != "") || + ((target.Mode == ModeSelfHosted || target.Mode == ModeManaged) && target.EnvironmentID == "") { + return Observation{}, errors.New("Runtime observation resolver returned mismatched Environment identity") + } if target.Mode == ModeNone || target.Mode == ModeSelfHosted { - return Observation{Target: target, Status: StatusUnsupported, Reason: "runtime_mode_not_observable"}, nil + return Observation{Target: target, Status: StatusUnsupported, Reason: "runtime_mode_not_observable", ResolvedAt: resolvedAt}, nil } if target.Mode != ModeManaged || target.Instance.AllocationID == "" || target.Instance.ProviderKey == "" { return Observation{}, errors.New("invalid managed Runtime observation target") } + if !target.Instance.AllocationCreatedAt.IsZero() && + (target.Instance.AllocationCreatedAt.Unix() < 0 || target.Instance.AllocationCreatedAt.After(resolvedAt)) { + return Observation{}, errors.New("invalid managed Runtime allocation creation time") + } + switch target.Instance.AllocationState { + case "creating": + return Observation{Target: target, Status: StatusUnavailable, Reason: "allocation_pending", ResolvedAt: resolvedAt}, nil + case "cleanup_pending", "released": + return Observation{Target: target, Status: StatusUnavailable, Reason: "runtime_not_running", ResolvedAt: resolvedAt}, nil + case "running": + default: + return Observation{}, errors.New("invalid managed Runtime allocation state") + } source, ok := s.sources[target.Instance.ProviderKey] if !ok { - return Observation{Target: target, Status: StatusUnavailable, Reason: "runtime_source_unavailable"}, nil + return Observation{Target: target, Status: StatusUnavailable, Reason: "source_not_configured", ResolvedAt: resolvedAt}, nil + } + providerType := "" + if typed, ok := source.(interface{ ObservationProviderType() string }); ok { + providerType = typed.ObservationProviderType() } sample, err := source.Observe(ctx, target) + if errors.Is(err, context.DeadlineExceeded) { + return Observation{Target: target, Status: StatusUnavailable, Reason: "sample_timeout", ProviderType: providerType, ResolvedAt: s.now()}, nil + } + if errors.Is(err, ErrNotRunning) { + return Observation{Target: target, Status: StatusUnavailable, Reason: "runtime_not_running", ProviderType: providerType, ResolvedAt: s.now()}, nil + } if errors.Is(err, ErrUnavailable) { - return Observation{Target: target, Status: StatusUnavailable, Reason: "runtime_sample_unavailable"}, nil + return Observation{Target: target, Status: StatusUnavailable, Reason: "sample_unavailable", ProviderType: providerType, ResolvedAt: s.now()}, nil } if err != nil { return Observation{}, fmt.Errorf("observe Runtime: %w", err) @@ -62,5 +98,5 @@ func (s *Service) ObserveSession(ctx context.Context, tenantID, sessionID string if err := sample.validate(s.now()); err != nil { return Observation{}, err } - return Observation{Target: target, Status: StatusObserved, Sample: &sample}, nil + return Observation{Target: target, Status: StatusObserved, Sample: &sample, ProviderType: providerType, ResolvedAt: s.now()}, nil } diff --git a/services/agents-api/internal/runtimeobs/service_test.go b/services/agents-api/internal/runtimeobs/service_test.go index bff6bf33d..6cb76682f 100644 --- a/services/agents-api/internal/runtimeobs/service_test.go +++ b/services/agents-api/internal/runtimeobs/service_test.go @@ -3,6 +3,7 @@ package runtimeobs import ( "context" "errors" + "math" "testing" "time" ) @@ -12,8 +13,15 @@ type fixedResolver struct { err error } -func (r fixedResolver) Resolve(context.Context, string, string) (Target, error) { - return r.target, r.err +func (r fixedResolver) Resolve(_ context.Context, tenant, session string) (Target, error) { + target := r.target + if target.TenantID == "" { + target.TenantID = tenant + } + if target.SessionID == "" { + target.SessionID = session + } + return target, r.err } type fixedSource struct { @@ -27,28 +35,48 @@ func (s *fixedSource) Observe(context.Context, Target) (Sample, error) { return s.sample, s.err } +type typedSource struct { + *fixedSource + providerType string +} + +func (s typedSource) ObservationProviderType() string { return s.providerType } + +type blockingSource struct{} + +func (blockingSource) Observe(ctx context.Context, _ Target) (Sample, error) { + <-ctx.Done() + return Sample{}, ctx.Err() +} + +func (blockingSource) ObservationProviderType() string { return "docker" } + func TestServiceDoesNotCallSourcesForUnsupportedModes(t *testing.T) { for _, mode := range []Mode{ModeNone, ModeSelfHosted} { source := &fixedSource{} - service, err := NewService(fixedResolver{target: Target{Mode: mode}}, map[string]Source{"provider": source}) + target := Target{Mode: mode} + if mode == ModeSelfHosted { + target.EnvironmentID = "environment" + } + service, err := NewService(fixedResolver{target: target}, map[string]Source{"provider": source}) if err != nil { t.Fatal(err) } observation, err := service.ObserveSession(t.Context(), "tenant", "session") - if err != nil || observation.Status != StatusUnsupported || source.calls != 0 { + if err != nil || observation.Status != StatusUnsupported || observation.Reason != "runtime_mode_not_observable" || source.calls != 0 { t.Fatalf("unsupported mode touched a source: %+v %v calls=%d", observation, err, source.calls) } } } func TestServicePreservesUnavailableAndObservedZero(t *testing.T) { - target := Target{Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider"}} + target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} service, err := NewService(fixedResolver{target: target}, nil) if err != nil { t.Fatal(err) } observation, err := service.ObserveSession(t.Context(), "tenant", "session") - if err != nil || observation.Status != StatusUnavailable || observation.Sample != nil { + if err != nil || observation.Status != StatusUnavailable || observation.Reason != "source_not_configured" || observation.Sample != nil { t.Fatalf("missing source was not unavailable: %+v %v", observation, err) } @@ -68,15 +96,20 @@ func TestServicePreservesUnavailableAndObservedZero(t *testing.T) { } func TestServiceMapsOnlyDeclaredUnavailability(t *testing.T) { - target := Target{Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider"}} + target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} for _, tc := range []struct { - err error - wantError bool + err error + wantReason string + wantError bool }{ - {err: ErrUnavailable}, + {err: ErrUnavailable, wantReason: "sample_unavailable"}, + {err: ErrNotRunning, wantReason: "runtime_not_running"}, + {err: context.DeadlineExceeded, wantReason: "sample_timeout"}, {err: errors.New("Docker permission denied"), wantError: true}, } { - service, err := NewService(fixedResolver{target: target}, map[string]Source{"provider": &fixedSource{err: tc.err}}) + service, err := NewService(fixedResolver{target: target}, map[string]Source{ + "provider": typedSource{fixedSource: &fixedSource{err: tc.err}, providerType: "docker"}, + }) if err != nil { t.Fatal(err) } @@ -84,8 +117,116 @@ func TestServiceMapsOnlyDeclaredUnavailability(t *testing.T) { if (err != nil) != tc.wantError { t.Fatalf("wrong error classification: %+v %v", observation, err) } - if !tc.wantError && observation.Status != StatusUnavailable { + if !tc.wantError && (observation.Status != StatusUnavailable || observation.Reason != tc.wantReason || observation.ProviderType != "docker") { t.Fatalf("declared unavailability was not mapped: %+v", observation) } } } + +func TestServiceMapsAnActualSourceDeadlineWithoutLeakingIt(t *testing.T) { + target := Target{ + EnvironmentID: "environment", Mode: ModeManaged, + Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}, + } + service, err := NewService(fixedResolver{target: target}, map[string]Source{"provider": blockingSource{}}) + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithTimeout(t.Context(), 10*time.Millisecond) + defer cancel() + observation, err := service.ObserveSession(ctx, "tenant", "session") + if err != nil || observation.Status != StatusUnavailable || observation.Reason != "sample_timeout" || observation.ProviderType != "docker" { + t.Fatalf("source deadline was not safely classified: %+v %v", observation, err) + } +} + +func TestServiceClassifiesResolverAndTerminalAllocationUnavailability(t *testing.T) { + now := time.Date(2026, 9, 22, 1, 0, 0, 0, time.UTC) + service, err := NewService(fixedResolver{target: Target{SessionID: "session", EnvironmentID: "environment", Mode: ModeManaged}, err: ErrUnavailable}, nil) + if err != nil { + t.Fatal(err) + } + service.now = func() time.Time { return now } + observation, err := service.ObserveSession(t.Context(), "tenant", "session") + if err != nil || observation.Status != StatusUnavailable || observation.Reason != "allocation_pending" || !observation.ResolvedAt.Equal(now) { + t.Fatalf("pending allocation was not classified: %+v %v", observation, err) + } + + source := &fixedSource{} + target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "creating"}} + service, err = NewService(fixedResolver{target: target}, map[string]Source{"provider": source}) + if err != nil { + t.Fatal(err) + } + observation, err = service.ObserveSession(t.Context(), "tenant", "session") + if err != nil || observation.Status != StatusUnavailable || observation.Reason != "allocation_pending" || source.calls != 0 { + t.Fatalf("creating allocation reached its provider: %+v %v calls=%d", observation, err, source.calls) + } + + target = Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "released"}} + service, err = NewService(fixedResolver{target: target}, map[string]Source{"provider": &fixedSource{}}) + if err != nil { + t.Fatal(err) + } + observation, err = service.ObserveSession(t.Context(), "tenant", "session") + if err != nil || observation.Status != StatusUnavailable || observation.Reason != "runtime_not_running" { + t.Fatalf("released allocation was not classified: %+v %v", observation, err) + } + + source = &fixedSource{} + target = Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{ + AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running", AllocationCreatedAt: time.Now().Add(time.Hour), + }} + service, err = NewService(fixedResolver{target: target}, map[string]Source{"provider": source}) + if err != nil { + t.Fatal(err) + } + if _, err := service.ObserveSession(t.Context(), "tenant", "session"); err == nil || source.calls != 0 { + t.Fatalf("future allocation creation reached its provider: %v calls=%d", err, source.calls) + } +} + +func TestServiceRejectsUnsafeProviderSamples(t *testing.T) { + target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} + now := time.Date(2026, 9, 22, 1, 0, 0, 0, time.UTC) + preEpoch := time.Unix(-1, 0).UTC() + tooLarge := uint64(1 << 53) + for _, sample := range []Sample{ + {ObservedAt: preEpoch}, + {ObservedAt: now, StartedAt: &preEpoch}, + {ObservedAt: now, CPUUsageSecondsTotal: float64Pointer(-1)}, + {ObservedAt: now, CPUUsageSecondsTotal: float64Pointer(math.NaN())}, + {ObservedAt: now, CPUUsageSecondsTotal: float64Pointer(math.Inf(1))}, + {ObservedAt: now, CPUCapacityCores: float64Pointer(0)}, + {ObservedAt: now, MemoryUsageBytes: &tooLarge}, + {ObservedAt: now, MemoryLimitBytes: &tooLarge}, + } { + service, err := NewService(fixedResolver{target: target}, map[string]Source{"provider": &fixedSource{sample: sample}}) + if err != nil { + t.Fatal(err) + } + service.now = func() time.Time { return now } + if _, err := service.ObserveSession(t.Context(), "tenant", "session"); err == nil { + t.Fatalf("unsafe sample accepted: %+v", sample) + } + } +} + +func TestServiceRejectsMismatchedResolvedOwnership(t *testing.T) { + for _, target := range []Target{ + {TenantID: "other", SessionID: "session", Mode: ModeNone}, + {TenantID: "tenant", SessionID: "other", Mode: ModeNone}, + {TenantID: "tenant", SessionID: "session", EnvironmentID: "unexpected", Mode: ModeNone}, + {TenantID: "tenant", SessionID: "session", Mode: ModeSelfHosted}, + } { + service, err := NewService(fixedResolver{target: target}, nil) + if err != nil { + t.Fatal(err) + } + if _, err := service.ObserveSession(t.Context(), "tenant", "session"); err == nil { + t.Fatalf("mismatched ownership accepted: %+v", target) + } + } +} + +func float64Pointer(value float64) *float64 { return &value } diff --git a/services/agents-api/internal/sandbox/docker/resources.go b/services/agents-api/internal/sandbox/docker/resources.go index 2f444c961..f0f057de0 100644 --- a/services/agents-api/internal/sandbox/docker/resources.go +++ b/services/agents-api/internal/sandbox/docker/resources.go @@ -15,6 +15,8 @@ import ( var _ runtimeobs.Source = (*Provider)(nil) +func (*Provider) ObservationProviderType() string { return "docker" } + // Observe is read-only. Inspect verifies allocation ownership before Docker // statistics are requested; it never renews or changes the container. func (p *Provider) Observe(ctx context.Context, target runtimeobs.Target) (runtimeobs.Sample, error) { @@ -27,13 +29,13 @@ func (p *Provider) Observe(ctx context.Context, target runtimeobs.Target) (runti reference := sandbox.Reference{TenantID: target.TenantID, EnvironmentID: target.EnvironmentID, AllocationID: target.Instance.AllocationID} inspected, err := p.inspect(ctx, reference) if errors.Is(err, sandbox.ErrNotFound) { - return runtimeobs.Sample{}, runtimeobs.ErrUnavailable + return runtimeobs.Sample{}, runtimeobs.ErrNotRunning } if err != nil { return runtimeobs.Sample{}, err } if inspected.Container.State == nil || !inspected.Container.State.Running { - return runtimeobs.Sample{}, runtimeobs.ErrUnavailable + return runtimeobs.Sample{}, runtimeobs.ErrNotRunning } result, err := p.client.ContainerStats(ctx, inspected.Container.ID, client.ContainerStatsOptions{Stream: false, IncludePreviousSample: false}) if err != nil { diff --git a/services/agents-api/internal/sandbox/docker/resources_test.go b/services/agents-api/internal/sandbox/docker/resources_test.go index ea4db1096..8a21dbff6 100644 --- a/services/agents-api/internal/sandbox/docker/resources_test.go +++ b/services/agents-api/internal/sandbox/docker/resources_test.go @@ -2,6 +2,7 @@ package docker import ( "encoding/json" + "errors" "net/http" "net/http/httptest" "strings" @@ -24,13 +25,14 @@ func TestObserveVerifiesOwnershipThenReadsOneShotStats(t *testing.T) { started := observed.Add(-time.Minute) statsRead := false omitMeasurements := false + running := true server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.Header().Set("Content-Type", "application/json") switch { case r.Method == http.MethodGet && strings.HasSuffix(r.URL.Path, "/json"): _ = json.NewEncoder(w).Encode(map[string]any{ "Id": "container-id", - "State": map[string]any{"Status": "running", "Running": true, "StartedAt": started.Format(time.RFC3339Nano)}, + "State": map[string]any{"Status": "running", "Running": running, "StartedAt": started.Format(time.RFC3339Nano)}, "HostConfig": map[string]any{"NanoCpus": 2_000_000_000, "Memory": 2048}, "Config": map[string]any{"Labels": map[string]string{ labelPrefix + "installation": installationID, @@ -76,6 +78,12 @@ func TestObserveVerifiesOwnershipThenReadsOneShotStats(t *testing.T) { if err != nil || missing.CPUUsageSecondsTotal != nil || missing.MemoryUsageBytes != nil || missing.CPUCapacityCores == nil || missing.MemoryLimitBytes == nil { t.Fatalf("missing Docker measurements became zero: %+v %v", missing, err) } + running = false + statsRead = false + if _, err := p.Observe(t.Context(), target); !errors.Is(err, runtimeobs.ErrNotRunning) || statsRead { + t.Fatalf("stopped Runtime was not classified before stats: %v stats=%v", err, statsRead) + } + running = true foreign := target foreign.Instance.ProviderKey = uuid.NewString() statsRead = false