diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index db71ce333..de8147bb5 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -859,8 +859,14 @@ native Run; they do not request model credentials or mutate the workspace. Cance retains the existing Turn/reservation semantics. Do not introduce a second queue. Publish metadata in the same transaction as Turn completion. Failed/cancelled Turns discard private objects, and Session deletion removes both private and published -copies. Reuse the source-file snapshot reader pattern and common content response; -artifact deletion does not alter workspace files. Hosted execution requires the +copies. The exporter skips output symlinks by their `lstat` type without following, +opening or resolving them; hard links, other special files, device crossings and +concurrent changes still reject the capture. In that completion transaction, drop +staged paths whose sha256 equals the newest remaining published Artifact for the +path in the Session, so later Turns publish only new, changed or no-longer-published +paths and never modify existing Artifacts. Reuse the source-file snapshot reader +pattern and common content response; artifact deletion does not alter workspace +files. Hosted execution requires the Runtime's bounded output-export capability and exact read-only preparation binding; capability advertisement alone does not qualify an operator's deployment. Exporter component checks do not establish public Artifact compatibility. diff --git a/contracts/agents-api/README.md b/contracts/agents-api/README.md index 57c90ac98..96fd68809 100644 --- a/contracts/agents-api/README.md +++ b/contracts/agents-api/README.md @@ -98,7 +98,7 @@ paths start at `/vaults`, not `/agents/vaults`. | sessions.events | create, stream | Text/cancel/function-result admission and live events; function-action state snapshots supported | | sessions.turns | retrieve, list | Implemented reads; lifecycle conformance still partial | | sessions.items | list | Partial Item variants | -| sessions.artifacts | retrieve, list, delete, content | Shared output capture and immutable stored reads/deletion on accepted Docker profiles and [qualified user-managed workflows](user-managed-runtime-v1.md) (prior Core-managed E2B evidence remains historical), including retained downloads after Runtime loss; exact upstream defaults/errors, unchanged-file republishing and cancellation-edge parity remain unverified | +| sessions.artifacts | retrieve, list, delete, content | Shared output capture and immutable stored reads/deletion on accepted Docker profiles and [qualified user-managed workflows](user-managed-runtime-v1.md) (prior Core-managed E2B evidence remains historical), including retained downloads after Runtime loss. [Aligned](official-semantics-alignment.md#artifact-capture-and-listing--september-23) output symlink skipping, unchanged-path non-republication, the list envelope and malformed filters; exact upstream defaults/errors, hard-link/special-file capture and cancellation-edge parity remain unverified | | sessions.subagents | retrieve, list | [Three-harness Docker reads, native lifecycle limits and real evidence](subagents.md); full multi-agent semantics remain partial | | sessions.subagents.items | list | Qualified own-child history reads; full Item variants and live child streaming remain partial | | sessions.subagents.turns | retrieve, list | Implemented; shared Session/child IDs | diff --git a/contracts/agents-api/official-semantics-alignment.md b/contracts/agents-api/official-semantics-alignment.md index 63428f2bc..1cd4cbeea 100644 --- a/contracts/agents-api/official-semantics-alignment.md +++ b/contracts/agents-api/official-semantics-alignment.md @@ -198,3 +198,75 @@ invalid bodies and queries, and replay U+0000 on every create/update family with a database digest proving no writes. The pinned-SDK acceptance scripts assert the new codes, params and messages. Independent real-Core acceptance is recorded separately by the coordinator. + +## Artifact capture and listing — September 23 + +This batch aligns Session Artifact capture and listing with the first official +Artifact observations. Evidence comes from the hosted-environment campaign scan +recorded privately in `~/.parsar/remediation/20260923/campaign-scan-2/hosted-env/` +(`findings.json` HE-50..62, raw records under `official/` and `run1/`). The probe +used three owned Sessions and two tiny `gpt-6-astra` Turns; all three Sessions +were deleted. Official Turn 1 created regular, nested and empty outputs plus +`outputs/link.txt -> a.txt`; Turn 2 only wrote `outputs/c.txt` after one Artifact +was deleted. + +| Row | Case | Core behavior | +| --- | --- | --- | +| A1 | A symlink below `outputs/` at Turn completion: to a file or directory, dangling, or pointing outside the workspace (HE-51) | Skipped by its `lstat` type: never followed, opened or resolved, and no Artifact. Every regular file is still captured and the Turn completes. | +| A2 | Later Turns in the same Session (HE-52) | A path is published again only when it has no remaining published Artifact in the Session, or its bytes (sha256) differ from the newest remaining one. Unchanged paths keep their existing Artifact IDs. The first Turn is unchanged. | +| A3 | List envelope (HE-53) | `object: list`, `data`, `first_id`, `last_id`, `has_more`, with null first/last IDs on an empty page, like the Session, Turn and Item lists. Paging and cursors are unchanged. | +| A4 | Malformed `environment_id` filter (HE-56) | 200 with an empty page, as for another existing Environment. Session lookup still runs first, so foreign and missing Sessions remain 404; cursor and limit errors are unchanged. | + +Decisions: + +- The Rust export helper handles every link kind the same way. Official evidence + shows one relative link to a file; telling the other kinds apart would require + resolving the link, which the confinement rules forbid. A link still counts as + a directory entry, so creating or removing one during export is a concurrent + change. +- The republication decision runs in the Turn's terminal transaction, not in the + private capture transaction. Capture commits and releases the Session lock + before the Turn completes, so an Artifact deletion can commit in between. The + terminal transaction holds the Session lock that also orders Artifact deletion, + and only one Turn per Session can be active, so the decision sees exactly the + Artifacts that remain at completion. Unchanged staged rows are deleted and their + private large objects unlinked in that transaction; published rows are never + modified. +- "Newest" follows the producing Turn's database creation time, then its ID. + Publication time can come from the Runtime's reported completion and is not a + reliable order between Turns. +- Known difference from the batch plan's wording, accepted as a local decision: + the plan republishes a path whose newest Artifact was deleted, but Core compares + against the newest *remaining* published Artifact. Deletion is physical and + leaves no record, and adding one would need a schema change outside this batch. + Example: Turn 1 publishes `b.txt` as `bravo`, Turn 2 publishes `bravo-v2`, and + the Turn 2 Artifact is then deleted. A later Turn whose `b.txt` is `bravo-v2` + republishes it, because the remaining Turn 1 version differs. A later Turn whose + `b.txt` is `bravo` publishes nothing, because the remaining Turn 1 Artifact + already has those bytes. The official behavior for this case is unobserved. +- A malformed filter resolves to the never-assigned maximum UUID, as for + malformed path identifiers, so it matches nothing without a database text + comparison. An empty `environment_id=` still means no filter. + +Deferred and unchanged: a linked `outputs` root, hard links, FIFOs, sockets, +devices and device crossings still reject the whole capture and fail the Turn +with `artifact_capture_failed`; there is no official evidence for them yet. +Republication after changed bytes is inferred rather than observed, and the +deleted-newest case above is unobserved. The unknown `after` cursor (HE-57) +belongs to ERROR-PROTOCOL-001. Subagent lists keep their `data`/`has_more` +envelope until there is official Subagent evidence. Artifact IDs keep the Core +UUID format. Paths removed from the workspace keep their Artifacts. + +Rust tests cover every link kind, including absolute links to a secret outside +the workspace and a relative link to a workspace file outside `outputs/`; an +inotify watch proves no target is opened or read, with a positive control. They +also keep the hard-link, socket, FIFO, linked-root and concurrent-change +rejections. Real-PostgreSQL store tests cover new, unchanged, changed, +changed-back, deleted-then-unchanged and deleted-during-capture paths, a deletion +that holds the Session lock while Turn completion waits, Turn-ordered newest +versions with inverted publication times, Session scoping and private object +accounting. Handler and real-PostgreSQL HTTP tests +cover the envelope, other, foreign and malformed filters, and foreign or missing +Sessions. The pinned-SDK and raw HTTP verifier used by live acceptance runs +against PostgreSQL across three Turns. Real Core, daemon and model acceptance is +recorded separately by the coordinator. diff --git a/contracts/agents-api/openapi.yaml b/contracts/agents-api/openapi.yaml index 7d5aabee9..8da7d07ce 100644 --- a/contracts/agents-api/openapi.yaml +++ b/contracts/agents-api/openapi.yaml @@ -1374,11 +1374,22 @@ definitions: items: $ref: '#/definitions/v1.SessionArtifact' type: array + first_id: + type: string + x-nullable: true has_more: type: boolean + last_id: + type: string + x-nullable: true + object: + enum: + - list + type: string required: - data - has_more + - object type: object v1.SessionDeleted: properties: @@ -3265,8 +3276,10 @@ paths: /agents/sessions/{session_id}/artifacts: get: description: Lists published outputs independently of Environment availability. - Sorting uses publication time and ID. The local default page size is 20; exact - upstream defaults and error parity remain unverified. + Sorting uses publication time and ID. A later Turn publishes a path again + only when it is new, its bytes changed, or no Artifact remains for it. A malformed + environment_id matches nothing. The local default page size is 20; exact upstream + defaults and error parity remain unverified. parameters: - description: agents=v1 in: header @@ -3278,7 +3291,8 @@ paths: name: session_id required: true type: string - - description: Producing Environment ID + - description: Producing Environment ID; an unknown or malformed ID returns + an empty page in: query name: environment_id type: string diff --git a/contracts/agents-api/operation-evidence.md b/contracts/agents-api/operation-evidence.md index bcd32745b..40f9d0579 100644 --- a/contracts/agents-api/operation-evidence.md +++ b/contracts/agents-api/operation-evidence.md @@ -38,6 +38,7 @@ Repository paths below are relative to the inspected worktree; private evidence | W | `~/.parsar/remediation/20260923/file-resource-semantics/official-files/` and `official-skills/`: 56 owned resource requests, no Sessions/models; [qualified observations and limitations](file-resource-semantics.md). | | X | [Validation error fields](official-semantics-alignment.md#validation-error-fields--september-23); private `~/.parsar/remediation/20260923/campaign-scan-1/{vaults-agents,sessions,skills-files-templates}/findings.json` VA-07/08/09/10, SES-28, SFT-20 and the September 22 Session `metadata.a` null observation. Metadata/name field errors, U+0000 local limit, malformed path IDs and Template network codes; Go handler and real-PostgreSQL route/no-write tests, no model execution. | | L | [List query tolerance](list-query-semantics.md#list-query-tolerance--september-23-2026); private `~/.parsar/remediation/20260923/campaign-scan-1/{vaults-agents,sessions,skills-files-templates}/findings.json`: owned-collection unknown/repeated keys, limit bounds, Vault status union and Files empty purpose, plus unknown keys on a deleted Vault read and Agent delete. Rows A1–D2 of that section; no model execution. | +| Y | [Artifact capture and listing](official-semantics-alignment.md#artifact-capture-and-listing--september-23); private `~/.parsar/remediation/20260923/campaign-scan-2/hosted-env/findings.json` HE-50..62 with raw records under `official/` (labels `al01`–`al09`, `ar01`–`ar04`, `ac01`–`ac05`, `ad01`/`ad02`): three owned Sessions and two tiny Turns, all deleted; first official Artifact observations. Rows A1–A4 of that section: symlink skip, republication, list envelope and malformed filter. Rust link tests, real-PostgreSQL store/HTTP and pinned-SDK tests without a model; live acceptance is recorded with the batch. | ## Per-operation evidence matrix @@ -60,10 +61,10 @@ Paths in the appendix include `/v1`. SDK names here omit `client.`. `P` means pa | 13 | beta.agents.sessions.turns.retrieve | P: persisted root/child Turn identity; malformed ID equals missing | H/S contain Turn list payloads; no isolated positive retrieve raw request identified in this set; X SES-28 official `turn_` 404 | Recorded H/B scoped Turn recovery; C history uses list | Distinguish list-shape evidence from retrieve wire qualification; full lifecycle/usage | | 14 | beta.agents.sessions.turns.list | P: full envelope, ordered root+child history; limit outside 1–100 rejects with the Beta code | S `turns-final/empty-page/limit-high/order-empty.json`; H both directions; L SES-14/15/16/17; X SES-28 | C Live paging; H/B real child/root identity recorded; L DB tenant A/B | All interleavings, same-timestamp paging, interim/failed usage | | 15 | beta.agents.sessions.items.list | P: scoped root Items, full envelope; limit 0/above 100 clamp within the pinned 1–100 page | S `items-final/empty-page/limit-high.json`; H both directions; L SES-12/13/15/16/17 | C Live; H/T recorded content/coordination/result variants; L DB tenant A/B | Full Item union. Newer turn_id filter excluded from pin | -| 16 | beta.agents.sessions.artifacts.retrieve | P: immutable captured metadata | None located | D recorded live Docker/user-managed workspace output | Exact hosted metadata/default/error and capture-edge parity | -| 17 | beta.agents.sessions.artifacts.list | P: scoped stored list | None located | D recorded live output enumeration | Paging during capture/delete; unchanged-file republishing | -| 18 | beta.agents.sessions.artifacts.delete | P: stored deletion | None located | D recorded workflow summary; per-case remote evidence not re-read | Repeated/in-flight deletion and physical retention parity | -| 19 | beta.agents.sessions.artifacts.content | P: immutable download after Runtime loss | None located | D recorded live retained download | Range/content headers, cancellation-edge capture and partial-transfer parity | +| 16 | beta.agents.sessions.artifacts.retrieve | P: immutable captured metadata; unchanged outputs keep their Artifact ID across Turns | Y HE-58/59 `ar01-retrieve`, missing and `ar03-cross-session` 404 (fields match; official `artifact_` IDs vs Core UUIDs) | D recorded live Docker/user-managed workspace output; Y DB unchanged-path ID and metadata stability | Error messages differ; hard-link/special-file capture edges unobserved | +| 17 | beta.agents.sessions.artifacts.list | P: scoped stored list, common list envelope; later Turns publish only new, changed or no-remaining-Artifact paths; symlinks skipped; malformed `environment_id` gives an empty page | Y HE-50..56 `al01`–`al09`: Turn 1 capture with a skipped link, Turn 2 republication, envelope, order/cursor/limits, other and malformed filters | D recorded live output enumeration; Y Rust link tests, DB republication, HTTP envelope/filter/tenant and pinned-SDK checks | Unknown `after` (HE-57, ERROR-PROTOCOL-001); changed-bytes republication inferred; deleted-newest edge and paging during capture/delete unobserved | +| 18 | beta.agents.sessions.artifacts.delete | P: stored deletion; the next Turn republishes a deleted path | Y HE-61 `ad01-delete`, `ad02-delete-repeat` 404, reads after delete 404; HE-52 deleted path republished by Turn 2 | D recorded workflow summary; Y DB deleted-then-unchanged and deleted-during-capture republication | In-flight deletion and physical retention parity | +| 19 | beta.agents.sessions.artifacts.content | P: immutable download after Runtime loss | Y HE-60/62 `ac01-content`, `ac02-content-range` (Range ignored, full 200), `ac05` 404 after Session deletion | D recorded live retained download; Y DB earlier versions keep their bytes | Core's extra Content-Disposition; cancellation-edge capture, partial transfer and content after Environment expiry unobserved officially | | 20 | beta.agents.sessions.subagents.retrieve | P: owned durable child identity/lifecycle | None located | B recorded Live six-read matrix | Full lifecycle/multi-agent parity; native close differences | | 21 | beta.agents.sessions.subagents.list | P: direct/nested/closed child records | None located | B recorded Live scopes/pages/continuation | Publication timing/parent propagation and unsupported native nesting | | 22 | beta.agents.sessions.subagents.items.list | P: child-owned history only | None located | B recorded Live child separation/recovery | Full child Item union; continuous child progress not qualified | @@ -110,7 +111,7 @@ Paths in the appendix include `/v1`. SDK names here omit `client.`. `P` means pa 2. **Session differences:** The Session admission batch removes idle `none` creation and empty metadata update. Local durable creation idempotency remains an explicit difference. Whitespace-only input succeeds officially but is rejected by the existing Core message validator; this newly observed difference is queued separately. Session agent updates, newer Environment shapes and root Item turn_id are baseline-upgrade questions. 3. **Template/Skill composition:** shared env/files/setup/packages selection is covered by the composition batch; template-reference null network/capability lists are covered by the null-selection batch. Official derived capability-directory projection remains different. Skill content/default metadata are covered by file-resource-semantics.md; sole-version deletion, visibility and broader numbering/error behavior remain unverified. 4. **Execution coverage:** use T's qualified matrix, not a blanket missing-image/structured-output claim. MiniMax functions/service MCP, optional tool combinations, unsupported images/placements and broader native lifecycle are explicit restrictions. PTC omission retains approved native behavior; Claude/MiniMax public Usage remains null; child settlement cadence/native close limits remain visible. No second executor/model loop or guessed counters are justified. -5. **Workspace and resources:** live Files bounds, symlink/path/cursor choices, artifact overwrite/republishing/headers/cancellation edges, full Environment metadata/lifecycle and Vault archive/in-flight-token semantics remain partial or unknown. Retired Core-managed E2B acceptance cannot qualify current user enrollment. +5. **Workspace and resources:** live Files bounds, symlink/path/cursor choices, artifact capture edges for hard links/special files and cancellation (Y aligns output symlinks, republication and the list envelope), full Environment metadata/lifecycle and Vault archive/in-flight-token semantics remain partial or unknown. Retired Core-managed E2B acceptance cannot qualify current user enrollment. ## Enumeration and wording mismatches diff --git a/contracts/agents-api/v1/session_artifacts.go b/contracts/agents-api/v1/session_artifacts.go index 78158a951..77f68951d 100644 --- a/contracts/agents-api/v1/session_artifacts.go +++ b/contracts/agents-api/v1/session_artifacts.go @@ -12,6 +12,9 @@ type SessionArtifact struct { } type SessionArtifactList struct { + Object string `json:"object" enums:"list" binding:"required"` + FirstID *string `json:"first_id" extensions:"x-nullable"` + LastID *string `json:"last_id" extensions:"x-nullable"` Data []SessionArtifact `json:"data" binding:"required"` HasMore bool `json:"has_more" binding:"required"` } diff --git a/packages/codex-executor/src/export.rs b/packages/codex-executor/src/export.rs index 5ba7df839..3e5de57b1 100644 --- a/packages/codex-executor/src/export.rs +++ b/packages/codex-executor/src/export.rs @@ -81,6 +81,9 @@ fn walk( )?; append(File::from(file), &path, device, budget, archive)?; } + // A link is recognized by its lstat type and skipped: it is never + // followed, opened or resolved, and publishes no Artifact. + FileType::Symlink => {} _ => return Err(invalid()), } } diff --git a/packages/codex-executor/src/export_tests.rs b/packages/codex-executor/src/export_tests.rs index 9fd81bfa4..cdbf5f874 100644 --- a/packages/codex-executor/src/export_tests.rs +++ b/packages/codex-executor/src/export_tests.rs @@ -1,6 +1,9 @@ use super::*; +use rustix::fs::{inotify, mkfifoat}; use std::collections::BTreeMap; +use std::ffi::CStr; use std::fs; +use std::mem::MaybeUninit; use std::os::unix::{fs::symlink, net::UnixListener}; #[test] @@ -46,10 +49,102 @@ fn absent_outputs_is_an_empty_archive() { ); } +/// Returns pending inotify events as (flags, entry name) without blocking. +fn drain(watcher: &OwnedFd) -> Vec<(inotify::ReadFlags, Option)> { + let mut buffer = [MaybeUninit::::uninit(); 8192]; + let mut reader = inotify::Reader::new(watcher, &mut buffer); + let mut events = Vec::new(); + loop { + match reader.next() { + Ok(event) => events.push(( + event.events(), + event.file_name().map(CStr::to_string_lossy).map(Into::into), + )), + Err(rustix::io::Errno::AGAIN) => return events, + Err(error) => panic!("inotify read: {error}"), + } + } +} + +#[test] +fn skips_symlinks_without_following_or_opening_their_targets() { + let _guard = crate::directory::TEST_LOCK.lock().unwrap(); + let root = tempfile::tempdir().unwrap(); + let other = tempfile::tempdir().unwrap(); + let secret = b"outside workspace secret marker"; + let private = b"workspace file outside outputs marker"; + fs::write(other.path().join("secret"), secret).unwrap(); + fs::create_dir(other.path().join("directory")).unwrap(); + fs::write(other.path().join("directory/secret"), secret).unwrap(); + fs::write(root.path().join("private.txt"), private).unwrap(); + fs::create_dir_all(root.path().join("outputs/sub")).unwrap(); + fs::write(root.path().join("outputs/a.txt"), "alpha").unwrap(); + fs::write(root.path().join("outputs/sub/b.txt"), "bravo").unwrap(); + fs::write(root.path().join("outputs/empty.txt"), []).unwrap(); + let tree = root.path().join("outputs"); + for (link, target) in [ + ("link.txt", Path::new("a.txt").to_path_buf()), + ("sub-link", "sub".into()), + ("sub/parent", "..".into()), + ("self", "self".into()), + ("dangling", "missing".into()), + ("private-link", "../private.txt".into()), + ("outside-file", other.path().join("secret")), + ("outside-directory", other.path().join("directory")), + ("sub/outside-root", other.path().into()), + ] { + symlink(target, tree.join(link)).unwrap(); + } + // Any open or read of a target through a followed link would be observed here. + let watcher = + inotify::init(inotify::CreateFlags::NONBLOCK | inotify::CreateFlags::CLOEXEC).unwrap(); + let flags = inotify::WatchFlags::OPEN | inotify::WatchFlags::ACCESS; + for watched in [ + other.path().to_path_buf(), + other.path().join("directory"), + root.path().join("private.txt"), + ] { + inotify::add_watch(&watcher, &watched, flags).unwrap(); + } + let mut bytes = Vec::new(); + outputs(root.path(), &mut bytes).unwrap(); + assert_eq!(drain(&watcher), [], "a link target was opened or read"); + let mut archive = tar::Archive::new(bytes.as_slice()); + let mut actual = BTreeMap::new(); + for entry in archive.entries().unwrap() { + let mut entry = entry.unwrap(); + assert!(entry.header().entry_type().is_file()); + let name = entry.path().unwrap().to_str().unwrap().to_owned(); + let mut data = Vec::new(); + entry.read_to_end(&mut data).unwrap(); + actual.insert(name, data); + } + assert_eq!( + actual, + BTreeMap::from([ + ("outputs/a.txt".into(), b"alpha".to_vec()), + ("outputs/empty.txt".into(), vec![]), + ("outputs/sub/b.txt".into(), b"bravo".to_vec()), + ]) + ); + for marker in [&secret[..], &private[..]] { + assert!(!bytes.windows(marker.len()).any(|part| part == marker)); + } + // Positive control: the watcher does report an actual target read. + assert_eq!(fs::read(other.path().join("secret")).unwrap(), secret); + assert!( + drain(&watcher) + .iter() + .any(|(flags, name)| flags.contains(inotify::ReadFlags::OPEN) + && name.as_deref() == Some("secret")), + "inotify control was not observed" + ); +} + #[test] -fn rejects_links_and_special_files_without_exposing_their_contents() { +fn rejects_root_link_hard_links_and_special_files_without_exposing_their_contents() { let _guard = crate::directory::TEST_LOCK.lock().unwrap(); - for kind in ["root-symlink", "file-symlink", "hardlink", "socket"] { + for kind in ["root-symlink", "hardlink", "socket", "fifo"] { let root = tempfile::tempdir().unwrap(); let other = tempfile::tempdir().unwrap(); let secret = b"private daemon credential marker"; @@ -58,18 +153,19 @@ fn rejects_links_and_special_files_without_exposing_their_contents() { symlink(other.path(), root.path().join("outputs")).unwrap(); } else { fs::create_dir(root.path().join("outputs")).unwrap(); + fs::write(root.path().join("outputs/regular"), "kept only on success").unwrap(); } let target = root.path().join("outputs/file"); let _socket = match kind { - "file-symlink" => { - symlink(other.path().join("secret"), target).unwrap(); - None - } "hardlink" => { fs::hard_link(other.path().join("secret"), target).unwrap(); None } "socket" => Some(UnixListener::bind(target).unwrap()), + "fifo" => { + mkfifoat(rustix::fs::CWD, &target, Mode::RUSR | Mode::WUSR).unwrap(); + None + } _ => None, }; let mut bytes = Vec::new(); @@ -105,17 +201,18 @@ fn rejects_a_changed_file_or_directory_during_export() { let _guard = crate::directory::TEST_LOCK.lock().unwrap(); struct Mutate<'a> { root: &'a Path, - directory: bool, + change: &'static str, done: bool, } impl Write for Mutate<'_> { fn write(&mut self, bytes: &[u8]) -> io::Result { if !self.done { self.done = true; - if self.directory { - fs::write(self.root.join("outputs/new"), "late")?; - } else { - fs::write(self.root.join("outputs/file"), "changed-length")?; + match self.change { + "file" => fs::write(self.root.join("outputs/file"), "changed-length")?, + "directory" => fs::write(self.root.join("outputs/new"), "late")?, + // Skipped links still take part in the directory listing check. + _ => symlink("file", self.root.join("outputs/late-link"))?, } } Ok(bytes.len()) @@ -124,7 +221,7 @@ fn rejects_a_changed_file_or_directory_during_export() { Ok(()) } } - for directory in [false, true] { + for change in ["file", "directory", "link"] { let root = tempfile::tempdir().unwrap(); fs::create_dir(root.path().join("outputs")).unwrap(); fs::write(root.path().join("outputs/file"), "initial").unwrap(); @@ -133,12 +230,12 @@ fn rejects_a_changed_file_or_directory_during_export() { root.path(), Mutate { root: root.path(), - directory, + change, done: false } ) .is_err(), - "directory={directory}" + "change={change}" ); } } diff --git a/services/agents-api/internal/api/session_artifacts.go b/services/agents-api/internal/api/session_artifacts.go index 5b049909b..2632b3881 100644 --- a/services/agents-api/internal/api/session_artifacts.go +++ b/services/agents-api/internal/api/session_artifacts.go @@ -31,13 +31,13 @@ func (h *Handler) artifactsReady(w http.ResponseWriter) bool { } // @Summary List immutable Session artifacts -// @Description Lists published outputs independently of Environment availability. Sorting uses publication time and ID. The local default page size is 20; exact upstream defaults and error parity remain unverified. +// @Description Lists published outputs independently of Environment availability. Sorting uses publication time and ID. A later Turn publishes a path again only when it is new, its bytes changed, or no Artifact remains for it. A malformed environment_id matches nothing. The local default page size is 20; exact upstream defaults and error parity remain unverified. // @Tags Artifacts // @Produce json // @Security BearerAuth // @Param OpenAI-Beta header string true "agents=v1" // @Param session_id path string true "Session ID" -// @Param environment_id query string false "Producing Environment ID" +// @Param environment_id query string false "Producing Environment ID; an unknown or malformed ID returns an empty page" // @Param after query string false "Last immutable artifact ID" // @Param limit query int false "Page size" minimum(1) maximum(100) default(20) // @Param order query string false "Publication order; omit for descending, explicit empty values are invalid" Enums(asc,desc) default(desc) @@ -57,11 +57,12 @@ func (h *Handler) listSessionArtifacts(w http.ResponseWriter, r *http.Request) { writeStoreError(w, r, err) return } - response := v1.SessionArtifactList{Data: make([]v1.SessionArtifact, 0, len(page.Artifacts)), HasMore: page.NextCursor != ""} + data := make([]v1.SessionArtifact, 0, len(page.Artifacts)) for _, artifact := range page.Artifacts { - response.Data = append(response.Data, artifactResponse(artifact)) + data = append(data, artifactResponse(artifact)) } - writeJSON(w, http.StatusOK, response) + first, last := listBounds(data, func(value v1.SessionArtifact) string { return value.ID }) + writeJSON(w, http.StatusOK, v1.SessionArtifactList{Object: "list", Data: data, HasMore: page.NextCursor != "", FirstID: first, LastID: last}) } // @Summary Retrieve immutable artifact metadata diff --git a/services/agents-api/internal/api/session_artifacts_test.go b/services/agents-api/internal/api/session_artifacts_test.go index 37b22cadb..7d2d35982 100644 --- a/services/agents-api/internal/api/session_artifacts_test.go +++ b/services/agents-api/internal/api/session_artifacts_test.go @@ -20,6 +20,7 @@ type artifactFixture struct { tenant, session, id, environment, cursor string limit int ascending bool + empty bool err error calls int } @@ -33,6 +34,9 @@ func (f *artifactFixture) GetSessionArtifact(_ context.Context, tenant, session, func (f *artifactFixture) ListSessionArtifacts(_ context.Context, tenant, session, environment, cursor string, limit int, ascending bool) (store.ArtifactPage, error) { f.calls++ f.tenant, f.session, f.environment, f.cursor, f.limit, f.ascending = tenant, session, environment, cursor, limit, ascending + if f.empty { + return store.ArtifactPage{Artifacts: []store.SessionArtifact{}}, f.err + } return store.ArtifactPage{Artifacts: []store.SessionArtifact{f.artifact}, NextCursor: f.artifact.ID}, f.err } @@ -78,6 +82,15 @@ func TestSessionArtifactRoutesAndPublicProjection(t *testing.T) { if w.Code != 200 || json.Unmarshal(w.Body.Bytes(), &page) != nil || !page.HasMore || len(page.Data) != 1 || f.environment != "environment" || f.cursor != "previous" || f.limit != 1 || !f.ascending { t.Fatalf("filtered page: %d %s %+v", w.Code, w.Body, f) } + // The list envelope matches the Session, Turn and Item lists (HE-53). + if page.Object != "list" || page.FirstID == nil || *page.FirstID != "artifact" || page.LastID == nil || *page.LastID != "artifact" { + t.Fatalf("list envelope: %s", w.Body) + } + f.empty = true + if w := request("GET", "?environment_id=not-a-uuid", "agents=v1"); w.Code != 200 || strings.TrimSpace(w.Body.String()) != `{"object":"list","first_id":null,"last_id":null,"data":[],"has_more":false}` || f.environment != "not-a-uuid" { + t.Fatalf("empty list envelope: %d %s", w.Code, w.Body) + } + f.empty = false w = request("GET", "/artifact/content", "agents=v1") if w.Code != 200 || !bytes.Equal(w.Body.Bytes(), []byte{0, 255, 1}) || w.Header().Get("Content-Type") != "application/octet-stream" || w.Header().Get("Content-Length") != "3" { t.Fatalf("content: %d %s %v", w.Code, w.Body, w.Header()) diff --git a/services/agents-api/internal/db/queries/session_artifacts.sql b/services/agents-api/internal/db/queries/session_artifacts.sql index 5d69cfa70..415e8d676 100644 --- a/services/agents-api/internal/db/queries/session_artifacts.sql +++ b/services/agents-api/internal/db/queries/session_artifacts.sql @@ -6,6 +6,25 @@ WHERE session_id = $1 AND id = $2 AND NOT artifact_capture_started; INSERT INTO session_artifacts (id, session_id, turn_id, environment_id, path, size_bytes, body_oid, sha256) VALUES ($1, $2, $3, $4, $5, $6, $7, $8); +-- name: DeleteUnchangedTurnArtifacts :exec +-- A staged path is not republished while the newest remaining published +-- Artifact for that path in the Session has the same bytes. Newest follows the +-- producing Turn's database creation time: Turns in a Session are serialized, +-- while publication time may come from the Runtime's reported completion. +WITH newest AS ( + SELECT DISTINCT ON (published.path) published.path, published.sha256 + FROM session_artifacts published + JOIN turns producer ON producer.session_id = published.session_id AND producer.id = published.turn_id + WHERE published.session_id = $1 AND published.created_at IS NOT NULL + ORDER BY published.path, producer.created_at DESC, producer.id DESC +), removed AS ( + DELETE FROM session_artifacts staged USING newest + WHERE staged.session_id = $1 AND staged.turn_id = $2 AND staged.created_at IS NULL + AND staged.path = newest.path AND staged.sha256 = newest.sha256 + RETURNING staged.body_oid +) +SELECT lo_unlink(body_oid) FROM removed; + -- name: PublishTurnArtifacts :exec UPDATE session_artifacts SET created_at = $3 WHERE session_id = $1 AND turn_id = $2 AND created_at IS NULL; diff --git a/services/agents-api/internal/db/sqlc/session_artifacts.sql.go b/services/agents-api/internal/db/sqlc/session_artifacts.sql.go index dbbb94d45..682972201 100644 --- a/services/agents-api/internal/db/sqlc/session_artifacts.sql.go +++ b/services/agents-api/internal/db/sqlc/session_artifacts.sql.go @@ -61,6 +61,36 @@ func (q *Queries) DeleteSessionArtifacts(ctx context.Context, sessionID pgtype.U return err } +const deleteUnchangedTurnArtifacts = `-- name: DeleteUnchangedTurnArtifacts :exec +WITH newest AS ( + SELECT DISTINCT ON (published.path) published.path, published.sha256 + FROM session_artifacts published + JOIN turns producer ON producer.session_id = published.session_id AND producer.id = published.turn_id + WHERE published.session_id = $1 AND published.created_at IS NOT NULL + ORDER BY published.path, producer.created_at DESC, producer.id DESC +), removed AS ( + DELETE FROM session_artifacts staged USING newest + WHERE staged.session_id = $1 AND staged.turn_id = $2 AND staged.created_at IS NULL + AND staged.path = newest.path AND staged.sha256 = newest.sha256 + RETURNING staged.body_oid +) +SELECT lo_unlink(body_oid) FROM removed +` + +type DeleteUnchangedTurnArtifactsParams struct { + SessionID pgtype.UUID `json:"session_id"` + TurnID pgtype.UUID `json:"turn_id"` +} + +// A staged path is not republished while the newest remaining published +// Artifact for that path in the Session has the same bytes. Newest follows the +// producing Turn's database creation time: Turns in a Session are serialized, +// while publication time may come from the Runtime's reported completion. +func (q *Queries) DeleteUnchangedTurnArtifacts(ctx context.Context, arg DeleteUnchangedTurnArtifactsParams) error { + _, err := q.db.Exec(ctx, deleteUnchangedTurnArtifacts, arg.SessionID, arg.TurnID) + return err +} + const deleteUnpublishedTurnArtifacts = `-- name: DeleteUnpublishedTurnArtifacts :exec WITH removed AS ( DELETE FROM session_artifacts WHERE session_id = $1 AND turn_id = $2 AND created_at IS NULL RETURNING body_oid diff --git a/services/agents-api/internal/store/artifact_lifecycle.go b/services/agents-api/internal/store/artifact_lifecycle.go index be0b004f1..2088c4083 100644 --- a/services/agents-api/internal/store/artifact_lifecycle.go +++ b/services/agents-api/internal/store/artifact_lifecycle.go @@ -45,8 +45,16 @@ func (s *Store) BeginTurnArtifactCapture(ctx context.Context, tenantID, sessionI }) } +// settleTurnArtifacts runs in the Turn's terminal transaction under the Session +// lock, which also orders Artifact deletion and allows one active Turn. The +// republication decision therefore sees exactly the Artifacts that remain when +// the Turn completes: new paths, changed bytes and paths whose newest Artifact +// was deleted are published; unchanged paths keep their existing Artifact IDs. func settleTurnArtifacts(ctx context.Context, q *sqlc.Queries, turn sqlc.Turn) error { if turn.Status == TurnCompleted { + if err := q.DeleteUnchangedTurnArtifacts(ctx, sqlc.DeleteUnchangedTurnArtifactsParams{SessionID: turn.SessionID, TurnID: turn.ID}); err != nil { + return err + } return q.PublishTurnArtifacts(ctx, sqlc.PublishTurnArtifactsParams{SessionID: turn.SessionID, TurnID: turn.ID, CreatedAt: turn.CompletedAt}) } return q.DeleteUnpublishedTurnArtifacts(ctx, sqlc.DeleteUnpublishedTurnArtifactsParams{SessionID: turn.SessionID, TurnID: turn.ID}) diff --git a/services/agents-api/internal/store/session_artifacts.go b/services/agents-api/internal/store/session_artifacts.go index 69d5df217..426b91aa4 100644 --- a/services/agents-api/internal/store/session_artifacts.go +++ b/services/agents-api/internal/store/session_artifacts.go @@ -50,11 +50,8 @@ func (s *Store) ListSessionArtifacts(ctx context.Context, tenantID, sessionID, e session, _ := parseID(sessionID) params := sqlc.ListSessionArtifactsParams{TenantID: tenant, SessionID: session, PageLimit: int32(limit + 1), Ascending: ascending, AfterID: pgtype.UUID{Valid: true}} if environmentID != "" { - var err error - params.EnvironmentID, err = parseID(environmentID) - if err != nil { - return ArtifactPage{}, err - } + // A malformed filter matches nothing, like another Environment's ID (HE-56). + params.EnvironmentID = parsePathID(environmentID) } if cursor != "" { // A malformed cursor remains an invalid request, unlike a path identifier. diff --git a/services/agents-api/internal/store/session_artifacts_public_test.go b/services/agents-api/internal/store/session_artifacts_public_test.go new file mode 100644 index 000000000..db23b5d0a --- /dev/null +++ b/services/agents-api/internal/store/session_artifacts_public_test.go @@ -0,0 +1,272 @@ +package store_test + +import ( + "archive/tar" + "bytes" + "context" + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "net/url" + "os" + "os/exec" + "reflect" + "strings" + "testing" + "time" + + "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/device" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/api" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/store" + "github.com/google/uuid" +) + +// hostedArtifactSession creates an openai_hosted Session without Turns. +func hostedArtifactSession(t *testing.T, s *store.Store, tenant, key string) (session, environment string) { + t.Helper() + created, err := s.CreateSession(t.Context(), tenant, store.CreateSessionInput{Creator: store.FixtureCreator(), Engine: "codex", IdempotencyKey: key, + Configuration: json.RawMessage(`{"agent":{"model":"artifact-model"},"environment":{"type":"openai_hosted","workspace_directory":"/workspace","capability_directories":[]}}`)}) + if err != nil || created.Environment == nil { + t.Fatal("fixture Session", err) + } + return created.ID, created.Environment.ID +} + +// completeArtifactTurn runs one Turn whose complete outputs tree is captured +// and settled as completed, and returns the Turn ID. +func completeArtifactTurn(t *testing.T, s *store.Store, tenant, session, environment, key string, outputs map[string]string) string { + t.Helper() + receipt, err := s.SubmitMessage(t.Context(), tenant, session, key, json.RawMessage(`{"input":[{"role":"user","content":[{"type":"input_text","text":"publish outputs"}]}]}`)) + if err != nil { + t.Fatal(err) + } + transition := func(from, to string) { + t.Helper() + if _, err := s.TransitionTurn(t.Context(), tenant, session, receipt.TurnID, store.TurnTransition{ExpectedStatus: from, Status: to}); err != nil { + t.Fatal(err) + } + } + transition(store.TurnQueued, store.TurnInProgress) + var archive bytes.Buffer + w := tar.NewWriter(&archive) + for name, body := range outputs { + if err := w.WriteHeader(&tar.Header{Name: "outputs/" + name, Mode: 0600, Size: int64(len(body)), Typeflag: tar.TypeReg}); err != nil { + t.Fatal(err) + } + if _, err := w.Write([]byte(body)); err != nil { + t.Fatal(err) + } + } + if err := w.Close(); err != nil { + t.Fatal(err) + } + if err := s.StageTurnArtifacts(t.Context(), tenant, session, receipt.TurnID, environment, &archive); err != nil { + t.Fatal(err) + } + transition(store.TurnInProgress, store.TurnCompleted) + return receipt.TurnID +} + +// artifactHTTPServer serves Artifact routes for an owner and a foreign tenant. +func artifactHTTPServer(t *testing.T, s *store.Store) (server *httptest.Server, owner, ownerTenant, foreign, foreignTenant string) { + t.Helper() + owner, foreign = uuid.NewString(), uuid.NewString() + ownerTenant, foreignTenant = uuid.NewString(), uuid.NewString() + auth, err := api.NewAuthenticator([]api.APIKey{ + {OrganizationID: "test-org", ProjectID: uuid.NewString(), SubjectKind: "service_account", SubjectID: "artifact-owner", TokenSHA256: device.HashCredential(owner), TenantID: ownerTenant}, + {OrganizationID: "test-org", ProjectID: uuid.NewString(), SubjectKind: "service_account", SubjectID: "artifact-foreign", TokenSHA256: device.HashCredential(foreign), TenantID: foreignTenant}, + }) + if err != nil { + t.Fatal(err) + } + h, err := api.NewHandler(s, auth, "codex", api.WithSessionArtifacts(s)) + if err != nil { + t.Fatal(err) + } + server = httptest.NewServer(h) + t.Cleanup(server.Close) + return server, owner, ownerTenant, foreign, foreignTenant +} + +// The Artifact list uses the common list envelope (HE-53), and a malformed +// environment_id filter matches nothing like another Environment's ID (HE-56), +// without weakening tenant or Session scoping. +func TestSessionArtifactListEnvelopeAndEnvironmentFilterPostgres(t *testing.T) { + s, _ := store.NewTestStore(t) + server, owner, ownerTenant, foreign, foreignTenant := artifactHTTPServer(t, s) + client := pathIDClient{t: t, server: server} + + session, environment := hostedArtifactSession(t, s, ownerTenant, "artifact-list") + completeArtifactTurn(t, s, ownerTenant, session, environment, "artifact-list-turn", map[string]string{"a.txt": "alpha", "b.txt": "bravo"}) + idle, otherEnvironment := hostedArtifactSession(t, s, ownerTenant, "artifact-idle") + foreignSession, foreignEnvironment := hostedArtifactSession(t, s, foreignTenant, "artifact-foreign") + completeArtifactTurn(t, s, foreignTenant, foreignSession, foreignEnvironment, "artifact-foreign-turn", map[string]string{"a.txt": "alpha"}) + + type envelope struct { + Object *string `json:"object"` + FirstID *string `json:"first_id"` + LastID *string `json:"last_id"` + Data []json.RawMessage `json:"data"` + HasMore *bool `json:"has_more"` + } + list := func(token, session string, query url.Values) (int, string, envelope) { + t.Helper() + status, raw := client.do(token, http.MethodGet, "/v1/agents/sessions/"+session+"/artifacts?"+query.Encode(), "", nil) + var page envelope + if status == http.StatusOK { + var keys map[string]json.RawMessage + if json.Unmarshal([]byte(raw), &keys) != nil || json.Unmarshal([]byte(raw), &page) != nil { + t.Fatalf("list body: %s", raw) + } + if len(keys) != 5 || keys["first_id"] == nil || keys["last_id"] == nil || page.Object == nil || *page.Object != "list" || page.Data == nil || page.HasMore == nil { + t.Fatalf("list envelope fields: %s", raw) + } + } + return status, raw, page + } + ids := func(page envelope) []string { + t.Helper() + out := make([]string, 0, len(page.Data)) + for _, item := range page.Data { + var value struct{ ID string } + if json.Unmarshal(item, &value) != nil || value.ID == "" { + t.Fatalf("artifact item: %s", item) + } + out = append(out, value.ID) + } + return out + } + const empty = `{"object":"list","first_id":null,"last_id":null,"data":[],"has_more":false}` + "\n" + + // First and last IDs bound every non-empty page; paging is unchanged. + status, raw, all := list(owner, session, url.Values{"order": {"asc"}}) + published := ids(all) + if status != http.StatusOK || len(published) != 2 || *all.FirstID != published[0] || *all.LastID != published[1] || *all.HasMore { + t.Fatalf("full page: %d %s", status, raw) + } + status, raw, first := list(owner, session, url.Values{"order": {"asc"}, "limit": {"1"}}) + if status != http.StatusOK || !reflect.DeepEqual(ids(first), published[:1]) || *first.FirstID != published[0] || *first.LastID != published[0] || !*first.HasMore { + t.Fatalf("first page: %d %s", status, raw) + } + status, raw, second := list(owner, session, url.Values{"order": {"asc"}, "limit": {"1"}, "after": {published[0]}}) + if status != http.StatusOK || !reflect.DeepEqual(ids(second), published[1:]) || *second.FirstID != published[1] || *second.LastID != published[1] || *second.HasMore { + t.Fatalf("second page: %d %s", status, raw) + } + status, raw, _ = list(owner, session, url.Values{"order": {"asc"}, "limit": {"1"}, "after": {published[1]}}) + if status != http.StatusOK || raw != empty { + t.Fatalf("page after the last Artifact: %d %s", status, raw) + } + if status, raw, filtered := list(owner, session, url.Values{"order": {"asc"}, "environment_id": {environment}}); status != http.StatusOK || !reflect.DeepEqual(ids(filtered), published) { + t.Fatalf("own Environment filter: %d %s", status, raw) + } + if status, raw, _ := list(owner, idle, nil); status != http.StatusOK || raw != empty { + t.Fatalf("Session without Artifacts: %d %s", status, raw) + } + + // Another existing, foreign, unknown or malformed Environment ID: an empty page. + for _, filter := range []string{otherEnvironment, foreignEnvironment, uuid.NewString(), "not-a-uuid", "env_" + environment, "\xff", "ffffffff-ffff-ffff-ffff-ffffffffffff"} { + for _, query := range []url.Values{{"environment_id": {filter}}, {"environment_id": {filter}, "after": {published[0]}, "limit": {"1"}}} { + if status, raw, _ := list(owner, session, query); status != http.StatusOK || raw != empty { + t.Errorf("environment_id=%q %v: %d %s", filter, query, status, raw) + } + } + } + // Cursor validation and page bounds keep their own errors under any filter. + for _, query := range []url.Values{{"environment_id": {"not-a-uuid"}, "after": {"not-a-uuid"}}, {"environment_id": {"not-a-uuid"}, "limit": {"0"}}} { + if status, raw, _ := list(owner, session, query); status != http.StatusBadRequest { + t.Errorf("%v: %d %s", query, status, raw) + } + } + if status, raw, _ := list(owner, session, url.Values{"environment_id": {"not-a-uuid"}, "after": {uuid.NewString()}}); status != http.StatusNotFound { + t.Errorf("unknown cursor with malformed filter: %d %s", status, raw) + } + + // Tenant and Session scoping precede the filter: foreign and missing Sessions stay 404. + missingStatus, missingBody := client.do(owner, http.MethodGet, "/v1/agents/sessions/"+uuid.NewString()+"/artifacts", "", nil) + if missingStatus != http.StatusNotFound { + t.Fatalf("missing Session: %d %s", missingStatus, missingBody) + } + for _, probe := range []struct{ token, session string }{{foreign, session}, {owner, foreignSession}, {owner, uuid.NewString()}, {owner, "not-a-uuid"}} { + for _, filter := range []string{"", "not-a-uuid", environment, foreignEnvironment} { + query := url.Values{} + if filter != "" { + query.Set("environment_id", filter) + } + if status, raw, _ := list(probe.token, probe.session, query); status != missingStatus || raw != missingBody { + t.Errorf("foreign or missing Session %s with filter %q: %d %s", probe.session, filter, status, raw) + } + } + } + // The foreign tenant still sees only its own Artifacts. + if status, raw, own := list(foreign, foreignSession, nil); status != http.StatusOK || len(own.Data) != 1 { + t.Fatalf("foreign tenant own list: %d %s", status, raw) + } +} + +// The pinned SDK and raw HTTP verifier used by live acceptance reads the +// republication result of three Turns, the list envelope and the filter. +func TestSessionArtifactsOfficialClientPostgres(t *testing.T) { + python := os.Getenv("PARSAR_OFFICIAL_SDK_PYTHON") + if python == "" { + t.Skip("pinned official Python SDK required") + } + s, _ := store.NewTestStore(t) + server, owner, ownerTenant, foreign, _ := artifactHTTPServer(t, s) + session, environment := hostedArtifactSession(t, s, ownerTenant, "artifact-sdk") + outputs := map[string]string{"a.txt": "alpha", "sub/b.txt": "bravo", "empty.txt": ""} + first := completeArtifactTurn(t, s, ownerTenant, session, environment, "artifact-sdk-1", outputs) + page, err := s.ListSessionArtifacts(t.Context(), ownerTenant, session, "", "", 100, true) + if err != nil { + t.Fatal(err) + } + for _, artifact := range page.Artifacts { + if artifact.Path == "/workspace/outputs/a.txt" { + if err := s.DeleteSessionArtifact(t.Context(), ownerTenant, session, artifact.ID); err != nil { + t.Fatal(err) + } + } + } + outputs["c.txt"] = "charlie" + second := completeArtifactTurn(t, s, ownerTenant, session, environment, "artifact-sdk-2", outputs) + outputs["sub/b.txt"] = "bravo-v2" + third := completeArtifactTurn(t, s, ownerTenant, session, environment, "artifact-sdk-3", outputs) + hex := func(files map[string]string) map[string]string { + out := make(map[string]string, len(files)) + for name, body := range files { + out["/workspace/outputs/"+name] = fmt.Sprintf("%x", body) + } + return out + } + settings, err := json.Marshal(map[string]any{"base": server.URL, "token": owner, "foreign": foreign, "session": session, "environment": environment, + "expected": map[string]map[string]string{ + first: hex(map[string]string{"sub/b.txt": "bravo", "empty.txt": ""}), + second: hex(map[string]string{"a.txt": "alpha", "c.txt": "charlie"}), + third: hex(map[string]string{"sub/b.txt": "bravo-v2"}), + }}) + if err != nil { + t.Fatal(err) + } + const driver = `import json, sys +sys.path.insert(0, "../../tests") +import httpx2 +from openai import DefaultHttpxClient, OpenAI +from official_session_artifacts import verify_session_artifacts +s = json.load(sys.stdin) +def client(key): + return OpenAI(api_key=key, base_url=s["base"] + "/v1", max_retries=0, + _strict_response_validation=True, http_client=DefaultHttpxClient(trust_env=False)) +expected = {turn: {path: bytes.fromhex(body) for path, body in files.items()} for turn, files in s["expected"].items()} +with httpx2.Client(trust_env=False, timeout=20) as http, client(s["token"]) as owner, client(s["foreign"]) as foreign: + items = verify_session_artifacts(owner, foreign, http, s["session"], s["environment"], expected) +print(json.dumps({"verified_artifacts": len(items)})) +` + ctx, cancel := context.WithTimeout(t.Context(), 2*time.Minute) + defer cancel() + command := exec.CommandContext(ctx, python, "-c", driver) + command.Stdin = bytes.NewReader(settings) + output, err := command.CombinedOutput() + if err != nil || !strings.Contains(string(output), `{"verified_artifacts": 5}`) { + t.Fatalf("official Artifact verification: %v %s", err, output) + } +} diff --git a/services/agents-api/internal/store/session_artifacts_test.go b/services/agents-api/internal/store/session_artifacts_test.go index 26f4ebe57..b9233a638 100644 --- a/services/agents-api/internal/store/session_artifacts_test.go +++ b/services/agents-api/internal/store/session_artifacts_test.go @@ -4,12 +4,17 @@ import ( "archive/tar" "bytes" "context" + "encoding/json" "errors" + "fmt" "io" "reflect" + "sort" + "strings" "testing" "time" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/db/sqlc" "github.com/google/uuid" ) @@ -267,3 +272,305 @@ func TestSessionArtifactTransferDoesNotBlockDeletionOrCancellation(t *testing.T) }) } } + +// startArtifactTurn admits another message and starts its Turn. +func startArtifactTurn(t *testing.T, s *Store, tenant, session, key string) string { + t.Helper() + input := submitMessage(t, s, tenant, session, key) + transition(t, s, tenant, session, input.TurnID, TurnQueued, TurnInProgress) + return input.TurnID +} + +// stageArtifactOutputs privately captures one complete outputs tree for a Turn. +func stageArtifactOutputs(t *testing.T, s *Store, tenant, session, environment, turn string, files map[string]string) { + t.Helper() + archive := make(map[string][]byte, len(files)) + for name, body := range files { + archive["outputs/"+name] = []byte(body) + } + if err := s.StageTurnArtifacts(t.Context(), tenant, session, turn, environment, bytes.NewReader(artifactArchive(t, archive))); err != nil { + t.Fatal(err) + } +} + +// publishedByTurn returns the Artifacts one Turn published, keyed by outputs-relative path. +func publishedByTurn(t *testing.T, s *Store, tenant, session, turn string) map[string]SessionArtifact { + t.Helper() + page, err := s.ListSessionArtifacts(t.Context(), tenant, session, "", "", 100, true) + if err != nil || page.NextCursor != "" { + t.Fatalf("list: %+v %v", page, err) + } + got := make(map[string]SessionArtifact) + for _, artifact := range page.Artifacts { + if artifact.TurnID == turn { + got[strings.TrimPrefix(artifact.Path, "/workspace/outputs/")] = artifact + } + } + return got +} + +func artifactBytes(t *testing.T, s *Store, tenant, session, id string) string { + t.Helper() + var body []byte + if err := s.ReadSessionArtifact(t.Context(), tenant, session, id, func(_ SessionArtifact, r io.Reader) error { + var err error + body, err = io.ReadAll(r) + return err + }); err != nil { + t.Fatal(err) + } + return string(body) +} + +func publishedPaths(published map[string]SessionArtifact) []string { + paths := make([]string, 0, len(published)) + for path := range published { + paths = append(paths, path) + } + sort.Strings(paths) + return paths +} + +// Later Turns publish a path only when it is new, its bytes differ from the +// newest remaining Artifact for that path, or no Artifact remains for it (HE-52). +func TestSessionArtifactsRepublishOnlyNewChangedOrDeletedPaths(t *testing.T) { + s, pool := testStore(t) + tenant, session, environment, first := artifactTurn(t, s, "openai_hosted") + before := sourceObjectCount(t, pool) + turnNumber := 1 + run := func(files map[string]string, want ...string) map[string]SessionArtifact { + t.Helper() + turn := first + if turnNumber > 1 { + turn = startArtifactTurn(t, s, tenant, session, fmt.Sprintf("artifact-turn-%d", turnNumber)) + } + turnNumber++ + stageArtifactOutputs(t, s, tenant, session, environment, turn, files) + transition(t, s, tenant, session, turn, TurnInProgress, TurnCompleted) + published := publishedByTurn(t, s, tenant, session, turn) + sort.Strings(want) + if got := publishedPaths(published); strings.Join(got, ",") != strings.Join(want, ",") { + t.Fatalf("Turn %d published %v, want %v", turnNumber-1, got, want) + } + for path, artifact := range published { + if body := artifactBytes(t, s, tenant, session, artifact.ID); body != files[path] { + t.Fatalf("Turn %d %s bytes = %q, want %q", turnNumber-1, path, body, files[path]) + } + } + return published + } + unchanged := func(artifacts ...SessionArtifact) { + t.Helper() + for _, artifact := range artifacts { + if got, err := s.GetSessionArtifact(t.Context(), tenant, session, artifact.ID); err != nil || got != artifact { + t.Fatalf("existing Artifact changed: %+v -> %+v %v", artifact, got, err) + } + } + } + objects := func() { + t.Helper() + var rows int + if err := pool.QueryRow(t.Context(), "SELECT count(*) FROM session_artifacts WHERE session_id = $1", session).Scan(&rows); err != nil { + t.Fatal(err) + } + if count := sourceObjectCount(t, pool); count != before+rows { + t.Fatalf("unpublished captures kept private objects: %d objects for %d Artifacts", count-before, rows) + } + } + + // The first Turn behaves as before: every regular output is published. + outputs := map[string]string{"a.txt": "alpha", "sub/b.txt": "bravo", "empty.txt": ""} + one := run(outputs, "a.txt", "sub/b.txt", "empty.txt") + if err := s.DeleteSessionArtifact(t.Context(), tenant, session, one["a.txt"].ID); err != nil { + t.Fatal(err) + } + // New c.txt and deleted-then-unchanged a.txt; unchanged paths keep their IDs. + outputs["c.txt"] = "charlie" + two := run(outputs, "a.txt", "c.txt") + unchanged(one["sub/b.txt"], one["empty.txt"]) + objects() + // Changed bytes publish a new version and leave the earlier one intact. + outputs["sub/b.txt"] = "bravo-v2" + three := run(outputs, "sub/b.txt") + unchanged(one["sub/b.txt"], one["empty.txt"], two["a.txt"], two["c.txt"]) + if body := artifactBytes(t, s, tenant, session, one["sub/b.txt"].ID); body != "bravo" { + t.Fatalf("earlier version changed: %q", body) + } + // Comparison uses the newest version, not any earlier one with equal bytes. + outputs["sub/b.txt"] = "bravo" + four := run(outputs, "sub/b.txt") + // Entirely unchanged outputs, and a removed workspace file, publish nothing. + delete(outputs, "c.txt") + run(outputs) + unchanged(one["sub/b.txt"], one["empty.txt"], two["a.txt"], two["c.txt"], three["sub/b.txt"], four["sub/b.txt"]) + objects() + // Deletion leaves no tombstone: the newest remaining version is the comparison base. + if err := s.DeleteSessionArtifact(t.Context(), tenant, session, four["sub/b.txt"].ID); err != nil { + t.Fatal(err) + } + outputs["sub/b.txt"] = "bravo-v2" + run(outputs) + outputs["sub/b.txt"] = "bravo" + run(outputs, "sub/b.txt") + + // A deletion committed after private capture but before completion is seen + // by the completion transaction, so the same Turn republishes the path. + turn := startArtifactTurn(t, s, tenant, session, "artifact-turn-delete-during-capture") + stageArtifactOutputs(t, s, tenant, session, environment, turn, outputs) + if err := s.DeleteSessionArtifact(t.Context(), tenant, session, two["a.txt"].ID); err != nil { + t.Fatal(err) + } + transition(t, s, tenant, session, turn, TurnInProgress, TurnCompleted) + if got := publishedPaths(publishedByTurn(t, s, tenant, session, turn)); !reflect.DeepEqual(got, []string{"a.txt"}) { + t.Fatalf("deletion during capture: published %v", got) + } + objects() + + // Another Session in the same tenant never compares against these Artifacts. + other, err := s.CreateSession(t.Context(), tenant, environmentInput("artifact-other", "openai_hosted", "/workspace")) + if err != nil { + t.Fatal(err) + } + otherEnvironment, err := s.GetSessionEnvironment(t.Context(), tenant, other.ID) + if err != nil { + t.Fatal(err) + } + otherTurn := startArtifactTurn(t, s, tenant, other.ID, "artifact-other-turn") + stageArtifactOutputs(t, s, tenant, other.ID, otherEnvironment.ID, otherTurn, outputs) + transition(t, s, tenant, other.ID, otherTurn, TurnInProgress, TurnCompleted) + if got := publishedPaths(publishedByTurn(t, s, tenant, other.ID, otherTurn)); !reflect.DeepEqual(got, []string{"a.txt", "empty.txt", "sub/b.txt"}) { + t.Fatalf("other Session first Turn published %v", got) + } + for _, id := range []string{session, other.ID} { + if err := s.DeleteSession(t.Context(), tenant, id); err != nil { + t.Fatal(err) + } + } + if count := sourceObjectCount(t, pool); count != before { + t.Fatalf("objects leaked: %d -> %d", before, count) + } +} + +// Publication time can come from the Runtime's reported completion and invert +// the order of Turns; the newest version for a path still follows Turn order. +func TestSessionArtifactsNewestVersionFollowsTurnOrder(t *testing.T) { + s, _ := testStore(t) + tenant := uuid.NewString() + created, err := s.CreateSession(t.Context(), tenant, environmentInput("artifact-order", "openai_hosted", "/workspace")) + if err != nil { + t.Fatal(err) + } + session := created.ID + env, err := s.GetSessionEnvironment(t.Context(), tenant, session) + if err != nil { + t.Fatal(err) + } + // Turn 1 reports a native completion one hour ahead, so its Artifact is + // published later than every following Turn's. + first := submitMessage(t, s, tenant, session, "artifact-order-1") + transition(t, s, tenant, session, first.TurnID, TurnQueued, TurnInProgress) + stageArtifactOutputs(t, s, tenant, session, env.ID, first.TurnID, map[string]string{"b.txt": "bravo"}) + future := time.Now().Add(time.Hour).UnixMilli() + if _, err := s.CompleteExecution(t.Context(), tenant, session, first.TurnID, TurnCompleted, json.RawMessage(fmt.Sprintf(`{"done":{"source_completed_at_ms":%d}}`, future)), "", first.Sequence); err != nil { + t.Fatal(err) + } + one := publishedByTurn(t, s, tenant, session, first.TurnID)["b.txt"] + run := func(key, body string) map[string]SessionArtifact { + t.Helper() + turn := startArtifactTurn(t, s, tenant, session, key) + stageArtifactOutputs(t, s, tenant, session, env.ID, turn, map[string]string{"b.txt": body}) + transition(t, s, tenant, session, turn, TurnInProgress, TurnCompleted) + return publishedByTurn(t, s, tenant, session, turn) + } + two := run("artifact-order-2", "bravo-v2")["b.txt"] + if two.ID == "" || !two.CreatedAt.Before(one.CreatedAt) { + t.Fatalf("fixture did not invert publication time: %+v %+v", one, two) + } + // Turn 2's version is the newest although Turn 1 was published later. + if got := run("artifact-order-3", "bravo-v2"); len(got) != 0 { + t.Fatalf("unchanged bytes of the newest Turn republished: %+v", got) + } + if got := publishedPaths(run("artifact-order-4", "bravo")); strings.Join(got, ",") != "b.txt" { + t.Fatalf("bytes of an older Turn's version were not republished: %v", got) + } +} + +// A deletion that holds the Session lock while Turn completion waits for it is +// seen by the completion transaction, which then republishes the path. +func TestSessionArtifactsCompletionWaitsForConcurrentDeletion(t *testing.T) { + s, pool := testStore(t) + tenant, session, environment, first := artifactTurn(t, s, "openai_hosted") + before := sourceObjectCount(t, pool) + stageArtifactOutputs(t, s, tenant, session, environment, first, map[string]string{"a.txt": "alpha"}) + transition(t, s, tenant, session, first, TurnInProgress, TurnCompleted) + newest := publishedByTurn(t, s, tenant, session, first)["a.txt"] + turn := startArtifactTurn(t, s, tenant, session, "artifact-concurrent-delete") + stageArtifactOutputs(t, s, tenant, session, environment, turn, map[string]string{"a.txt": "alpha"}) + + // Delete exactly as DeleteSessionArtifact does, but keep the transaction open. + tx, err := pool.Begin(t.Context()) + if err != nil { + t.Fatal(err) + } + defer tx.Rollback(context.Background()) + lookup, err := artifactLookup(tenant, session, newest.ID) + if err != nil { + t.Fatal(err) + } + q := s.queries.WithTx(tx) + if _, err := q.LockSession(t.Context(), sqlc.LockSessionParams{TenantID: lookup.TenantID, ID: lookup.SessionID}); err != nil { + t.Fatal(err) + } + oid, err := q.DeleteSessionArtifact(t.Context(), sqlc.DeleteSessionArtifactParams(lookup)) + if err != nil { + t.Fatal(err) + } + objects := tx.LargeObjects() + if err := objects.Unlink(t.Context(), oid.Uint32); err != nil { + t.Fatal(err) + } + done := make(chan error, 1) + go func() { + _, err := s.TransitionTurn(t.Context(), tenant, session, turn, TurnTransition{ExpectedStatus: TurnInProgress, Status: TurnCompleted}) + done <- err + }() + // Completion must be blocked on the Session lock before the deletion commits. + deadline := time.Now().Add(10 * time.Second) + for { + var waiting int + if err := pool.QueryRow(t.Context(), "SELECT count(*) FROM pg_stat_activity WHERE datname = current_database() AND wait_event_type = 'Lock'").Scan(&waiting); err != nil { + t.Fatal(err) + } + if waiting > 0 { + break + } + select { + case err := <-done: + t.Fatalf("completion did not wait for the Session lock: %v", err) + default: + } + if time.Now().After(deadline) { + t.Fatal("completion never waited for the Session lock") + } + time.Sleep(10 * time.Millisecond) + } + if status, err := s.GetTurn(t.Context(), tenant, session, turn); err != nil || status.Status != TurnInProgress { + t.Fatalf("Turn settled while the deletion held the lock: %+v %v", status, err) + } + if err := tx.Commit(t.Context()); err != nil { + t.Fatal(err) + } + if err := <-done; err != nil { + t.Fatal(err) + } + published := publishedByTurn(t, s, tenant, session, turn) + if got := publishedPaths(published); strings.Join(got, ",") != "a.txt" || artifactBytes(t, s, tenant, session, published["a.txt"].ID) != "alpha" { + t.Fatalf("concurrent deletion was not republished: %v", got) + } + if _, err := s.GetSessionArtifact(t.Context(), tenant, session, newest.ID); !errors.Is(err, ErrNotFound) { + t.Fatalf("deleted Artifact remains: %v", err) + } + if count := sourceObjectCount(t, pool); count != before+1 { + t.Fatalf("private objects: %d -> %d, want one published Artifact", before, count) + } +} diff --git a/services/agents-api/tests/official_session_artifacts.py b/services/agents-api/tests/official_session_artifacts.py index b86f25af0..037fe2b7f 100644 --- a/services/agents-api/tests/official_session_artifacts.py +++ b/services/agents-api/tests/official_session_artifacts.py @@ -1,7 +1,21 @@ -"""Pinned SDK and raw HTTP checks for known immutable Session output versions.""" +"""Pinned SDK and raw HTTP checks for known immutable Session output versions. + +`expected` maps each Turn ID to the outputs it published. A later Turn +publishes a path only when it is new, its bytes changed, or no Artifact +remains for it, so unchanged outputs stay under their earlier Turn. +""" from openai import NotFoundError +EMPTY_PAGE = {"object": "list", "data": [], "first_id": None, "last_id": None, "has_more": False} + + +def check_envelope(page): + """Assert the common list envelope used by Session, Turn and Item lists.""" + assert set(page) == {"object", "data", "first_id", "last_id", "has_more"} and page["object"] == "list" + ids = [item["id"] for item in page["data"]] + assert (page["first_id"], page["last_id"]) == ((ids[0], ids[-1]) if ids else (None, None)) + def verify_session_artifacts(client, foreign, http, session_id, environment_id, expected): resource = client.beta.agents.sessions.artifacts @@ -40,6 +54,7 @@ def verify_session_artifacts(client, foreign, http, session_id, environment_id, assert response.status_code == 200 page = response.json() assert isinstance(page["data"], list) and type(page["has_more"]) is bool + check_envelope(page) assert len(page["data"]) <= 2 seen.extend(item["id"] for item in page["data"]) if not page["has_more"]: @@ -52,6 +67,12 @@ def verify_session_artifacts(client, foreign, http, session_id, environment_id, assert [item.id for item in resource.list(session_id)] == list(reversed(ascending_ids)) assert [item.id for item in resource.list(session_id, after=None, environment_id=None, limit=None)] == list(reversed(ascending_ids)) + # Another, unknown or malformed Environment ID matches nothing. + for other_environment in ("00000000-0000-4000-8000-000000000000", "not-a-uuid"): + response = http.get(endpoint, headers=headers, params={"environment_id": other_environment}) + assert response.status_code == 200 and response.json() == EMPTY_PAGE + assert list(resource.list(session_id, environment_id=other_environment)) == [] + assert http.get(endpoint, headers={"Authorization": headers["Authorization"]}).status_code == 400 assert http.get(endpoint, headers={"OpenAI-Beta": "agents=v1"}).status_code == 401 for query in ("limit=0", "limit=101", "order=wrong", "environment_id=a&environment_id=b"): diff --git a/services/agents-api/tests/official_user_runtime.py b/services/agents-api/tests/official_user_runtime.py index cd7275c5e..e7cadccac 100644 --- a/services/agents-api/tests/official_user_runtime.py +++ b/services/agents-api/tests/official_user_runtime.py @@ -312,7 +312,7 @@ def isolation_proof(phase): isolation_proof("recovered") assert memory in answer(items), "Cold continuation lost native conversation history" assert runtime.read(prefix + "-published") == b"published\n", "Cold continuation repeated side effects" - artifacts[second.id] = outputs + # Unchanged outputs keep their first-Turn Artifacts; later Turns publish nothing new. check_artifacts() report["checks"].append("runtime_and_core_restart_preserve_history_and_outputs") third, items = run_turn("Use your native shell tool to run exactly `python3 " + prefix + "-hold.py` once and wait in the foreground. " @@ -332,7 +332,6 @@ def isolation_proof(phase): assert memory in answer(items) and runtime.read(prefix + "-resumed") == b"resumed\n", "Post-cancel continuation failed" assert runtime.read(starts) == b"started\n" and runtime.read(ticks) == stopped, "Continuation repeated cancelled effects" assert runtime.read(prefix + "-published") == b"published\n", "Continuation repeated publication" - artifacts[fourth.id] = outputs check_artifacts() report["checks"].append("public_cancel_stops_ticks_duplicate_cancel_and_continuation_preserve_effect_counts") report.update(passed=True, memory_sha256=hashlib.sha256(memory.encode()).hexdigest(),