Repository navigation
feat(storage): object-storage worker with S3/GCS/R2/local backends + multi-provider e2e harness - #91
Conversation
|
No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review info⚙️ Run configurationConfiguration used: Organization UI Review profile: CHILL Plan: Pro Run ID: 📒 Files selected for processing (2)
✅ Files skipped from review due to trivial changes (1)
🚧 Files skipped from review as they are similar to previous changes (1)
📝 WalkthroughWalkthroughAdds a new Rust storage worker crate with S3/GCS/R2/local backends, a Backend trait, RPC handlers (put/get/delete/presign), event trigger system and pollers, rustfs sidecar lifecycle, CLI/runtime wiring and manifest, README/configs, CI workflow for E2E, and extensive unit/integration/e2e tests plus a TypeScript harness and Docker fixtures. ChangesStorage worker end-to-end feature
Sequence Diagram(s)sequenceDiagram
participant Client
participant Worker
participant Backend
participant Poller
participant Engine
Client->>Worker: storage::* RPC
Worker->>Backend: put/get/delete/presign
Backend-->>Worker: Result/Error
Poller-->>Worker: Provider event
Worker->>Engine: trigger(functionId,payload,timeout)
Engine-->>Worker: ack/null/error
Estimated code review effort🎯 5 (Critical) | ⏱️ ~120 minutes Suggested reviewers
Poem
✨ Finishing Touches🧪 Generate unit tests (beta)
|
…e harness
New `storage/` worker that exposes a provider-agnostic object-storage
surface to the iii engine. Four RPC functions (putObject, getObject,
deleteObject, presignUrl) and two trigger types (object-created,
object-deleted) over four backends (S3, GCS, R2, local rustfs) behind
a single `Backend` trait, configured per-bucket. Worker, crate, binary,
and RPC namespace are all `storage` from day one.
Worker (production code)
- src/config.rs — YAML parsing for `providers:` + `buckets:` map.
Per-bucket credentials redacted in Debug; `${ENV}` expansion fails
the load on missing vars rather than silently producing empty
secrets. endpoint_url field on S3BucketConfig, R2BucketConfig,
GcsBucketConfig.
- src/error.rs — wire-stable `StorageError` (CONFIG_ERROR,
UNKNOWN_BUCKET, OBJECT_NOT_FOUND, BODY_TOO_LARGE, OBJECT_TOO_LARGE,
INVALID_BASE64, INVALID_PRESIGN_PARAMS, PRESIGN_UNSUPPORTED,
LOCAL_BACKEND_*, PROVIDER_ERROR, PROVIDER_AUTH_FAILED,
CF_QUEUE_AUTH_FAILED, TRIGGER_DISPATCH_TIMEOUT) plus BackendError →
StorageError mapping.
- src/backend/{mod,s3,gcs,r2,local,factory}.rs — `Backend` trait + four
concrete impls. R2 and local reuse the S3 backend with a custom
endpoint; GCS uses yoshidan's gcloud-storage 1.x (JSON API, native
auth). factory.rs threads endpoint_url to S3BackendOpts /
GcsBackendOpts / R2BackendOpts. R2 build emits tracing::warn when an
endpoint_url override is in effect (accidental production use is
visible in logs). GcsBackend::build assigns
ClientConfig.storage_endpoint when endpoint_url is set
(gcloud-storage 1.3.0 threads it into v1_endpoint and
v1_upload_endpoint for all object operations).
- src/handlers/{put,get,delete,presign}_object.rs — request/response
types, base64 framing, 10 MiB inline body cap, key validation,
consistent CONFIG_ERROR mapping.
- src/rustfs/{spawn,health}.rs — managed rustfs sidecar: `RUSTFS_BIN`
→ ./rustfs → `which` discovery, kernel-allocated port, ephemeral
credentials, S3-endpoint readiness probe, SIGTERM-then-SIGKILL
shutdown with logged outcome.
- src/triggers/normalize.rs — provider-specific event JSON →
ObjectEventNormalized. S3 (Records[]), GCS (attributes/data), R2
(action), rustfs (Records[] with EventName). Percent-decodes the
S3-spec object key once at the normalizer boundary so handlers
receive `harness/foo` rather than `harness%2Ffoo`; GCS uses raw
objectId in Pub/Sub attributes per its own spec.
- src/triggers/pollers/{sqs,pubsub,cf_queue,rustfs_webhook}.rs — one
poller per upstream source. SQS long-polls, Pub/Sub pulls in
batches, CF Queue uses the REST consume API (with startup auth
probe), rustfs binds a loopback HTTP receiver.
- src/triggers/{registry,dispatcher,handler}.rs — static topology:
one poller per upstream, fan-out to subscribers via TriggerRegistry
+ EngineDispatcher. ObjectCreated/Deleted handlers reject
registrations whose bucket has no configured notifications source.
- src/main.rs — wiring: parse config, optionally spawn rustfs (with
the loopback notify URL passed at spawn time via
RUSTFS_NOTIFY_WEBHOOK_ENABLE_iii / ENDPOINT_iii), build backends,
start pollers per notifications source, register 4 functions + 2
trigger types with iii-sdk.
Pre-landing review fixes folded in
- backend/gcs.rs uses CredentialsFile::new_from_file instead of
mutating GOOGLE_APPLICATION_CREDENTIALS (cross-bucket race +
subprocess leak).
- backend/{s3,gcs}.rs reject oversized objects on Content-Length
BEFORE buffering the body (memory-amplification DoS).
- GCS PUT switches to UploadType::Multipart when cache_control or
metadata is set so those fields actually land.
- GCS get pins download_object via if_generation_match to eliminate
the metadata/body race.
- s3.rs map_s3_error sources inner_code from
ProvideErrorMetadata::code() instead of parsing a Debug print.
Tests (in-tree)
- 67 unit tests (config, error mapping, backend mock, handlers,
normalize, trigger registry).
- tests/integration.rs — schema-presence regression for the 4 RPC
registrations.
- tests/e2e/{local,s3,gcs,r2}.rs — full round-trips. local runs
whenever rustfs is discoverable; the cloud three are #[ignore]'d
behind their respective env vars.
- Regression tests for percent-decoded triggers in each affected
normalizer.
E2E TypeScript harness (tests/e2e/workers/harness/)
- Provider type, BUCKETS map, forEachProvider helper in cases.ts.
- buildSchemaReset / buildFunctionCases accept the active provider
list at runtime; existing 5 RPC cases run once per provider.
- runner.ts probes each provider before running cases; cases
targeting a down provider report ERROR (env, not regression).
- report.json grows summary + by_provider rollup.
- worker.ts emits tri-state HARNESS_DONE sentinel (PASS|FAIL|ERROR).
- Empty-filter result is now an explicit FAIL, not silent PASS 0/0.
Edge / provider / concurrency cases
- Body shapes: empty, 1-byte, 5 MiB random binary, 0..255 byte range.
- Key shapes: spaces, plus, percent, equals, question, unicode,
trailing slash, deep nesting (50 segments).
- Error envelopes: OBJECT_NOT_FOUND on get, presign on missing key
succeeds, presign TTL floor, idempotent delete.
- Provider quirks: S3 presigned URL contains bucket name; R2 region
"auto" performs a real round-trip.
- Concurrency: 8-way concurrent overwrite asserts last-write-wins
consistency. Rapid-fire 10 sequential puts to distinct keys
asserts every event delivered.
Multi-provider runner (tests/e2e/run-tests.sh)
- --providers=local|all|<csv> flag (default local preserves today's
behavior).
- docker-compose.yml: MinIO in profile=cloud (serves both scratch-s3
and scratch-r2). --providers=local skips compose entirely.
- fixtures/minio-init.sh: idempotent bucket bootstrap via aws-cli in
a docker container; no host install required.
- config.all.yaml: four-bucket engine config (scratch-local, -s3, -r2).
- Pre-flight checks for docker, compose v2, and port availability.
Compose up/down lifecycle. Tri-state HARNESS_DONE sentinel and
exit codes 0/1/2/3 for pass/regression/env/infra. FATAL paths tail
engine + harness logs and per-container docker logs.
Script self-tests (tests/e2e/script-tests/run.sh)
Bash tests of run-tests.sh itself — unknown --providers exits 3,
unknown arg exits 2, -h prints help, port 49134 in use exits 3,
--filter zero-match yields explicit FAIL, SIGINT mid-run leaves no
rustfs/storage orphans. Independent from the harness.
CI (.github/workflows/storage-e2e.yml)
Two parallel jobs — script-tests (~3 min cold, <1 min warm) and
e2e --providers=all (~7 min cold, ~3 min warm with rust + node
caches). Triggered on PR paths touching storage/** and on manual
workflow_dispatch.
GCS descope (post-Phase-2)
The Phase 2 spike confirmed gcloud-storage 1.3.0 exposes
ClientConfig.storage_endpoint, but integration testing found the
crate eagerly authenticates against the real Google OAuth endpoint
at with_credentials().await time. With dummy SA JSON, Google
returns 400 invalid_grant before the worker registers any RPCs,
killing it. The crate has no anonymous mode (token_source_provider
is a required non-optional field; default is a no-op that errors
at request time). The endpoint_url field stays in production code;
the e2e harness drops fake-gcs-server, scratch-gcs, and the gcs
provider. Re-enabling requires either an upstream gcloud-storage
change or a swap to a different crate.
Operator-facing
- README.md, config.yaml.example — full config schema, per-provider
notes, error code table, and the SQS / Pub-Sub / CF Queue / rustfs
notification wiring instructions. README documents endpoint_url.
- iii.worker.yaml — manifest declaring the worker as a Rust binary.
Soft-ship caveat: R2's CF Queues consume REST API is the youngest of
the four notification sources; if pagination or auth-scope edge cases
surface in production, cf_queue.rs is the place to adjust. The
startup auth probe and CF_QUEUE_AUTH_FAILED error code give operators
a fast-fail signal.
Results
- cargo test: 70 passed, 3 ignored. cargo clippy clean.
- --providers=local: HARNESS_DONE: PASS 34/34
- --providers=all: HARNESS_DONE: PASS 84/84 (24 local, 25 s3,
25 r2, 10 unscoped triggers + errors)
- script-tests: SCRIPT_TESTS_DONE: PASS 15/15
Spec: docs/superpowers/specs/2026-05-06-iii-storage-e2e-hardening-design.md
Plan: docs/superpowers/plans/2026-05-06-iii-storage-e2e-hardening.md
(Spec/plan filenames retained as point-in-time artifacts of when the
worker was named iii-storage during planning; their content is
unchanged.)
…cratch-s3
Trigger-dispatch tests previously fanned out only against the local rustfs
backend; the scratch-s3 path through the worker's SQS poller had no e2e
coverage. The blocker was infrastructural — MinIO has no SQS notification
target, so we add a tiny Bun shim that translates MinIO bucket-notification
webhooks into messages on a co-located ElasticMQ queue, which the storage
worker then consumes via its existing SQS poller.
Compose now runs three new services in the cloud profile: elasticmq (SQS-
compatible broker), minio-sqs-bridge (webhook → SQS shim, idempotently
creates the queue on boot), and the MinIO container itself gains
MINIO_NOTIFY_WEBHOOK_* env so the webhook target is loaded at boot without
needing an admin restart from minio-init.sh. The init script binds
arn:minio:sqs::BRIDGE:webhook on scratch-s3 for put/delete events.
The harness runner now registers triggers per-provider, smoke-tests each
non-local provider's delivery path before running its trigger cases, and
ERROR-skips cases for providers whose probe fails (so a half-broken cloud
profile doesn't FAIL otherwise-healthy suites). cases-triggers.ts exports a
buildTriggerCases(providers) builder; the per-scenario timeout grows from
5s to 10s to absorb MinIO→bridge→ElasticMQ→long-poll latency.
run-tests.sh exports placeholder AWS_ACCESS_KEY_ID/SECRET/REGION when s3
is in --providers so the worker's SQS client can sign requests without
leaking the developer's real aws sso session into the test.
R2 trigger coverage is intentionally out of scope: no usable Cloudflare
Queue emulator exists; that path should be exercised via Rust integration
tests against EventDispatcher.
Verified: tsc --noEmit clean on harness and bridge; docker compose config
validates; script-tests/run.sh 15/15 PASS; live smoke test of the full
chain (aws s3 cp into scratch-s3 → SQS receive-message returns the
expected s3:ObjectCreated:Put record with Records[0].s3.{bucket,object}
shape that normalize.rs::extract_event consumes after decode_s3_key).
Add .dockerignore so a local `bun install` doesn't ship node_modules into the build context, and switch the bridge install to --frozen-lockfile + COPY bun.lock so CI builds resolve @aws-sdk/client-sqs to an exact pinned version instead of drifting on each run.
Three independent bugs prevented the probe — and therefore every s3
trigger case — from succeeding once the harness ran end-to-end.
1. Worker SQS client was hitting real AWS, not ElasticMQ.
aws-sdk-rust's SqsClient does not derive its endpoint from QueueUrl —
without an override it sends ReceiveMessage to
sqs.us-east-1.amazonaws.com, which rejected our placeholder creds with
InvalidClientTokenId (logged by the worker as opaque "service error").
Setting AWS_ENDPOINT_URL_SQS in run-tests.sh routes the SQS client to
the local ElasticMQ container.
2. MinIO emits eventName as "s3:ObjectCreated:Put" but the storage worker's
normalize_s3 (storage/src/triggers/normalize.rs:64-69) checks
starts_with("ObjectCreated:") — AWS-spec correct, but the prefix mismatch
meant every record hit the silent `continue` branch and got acked without
dispatch. The bridge now strips the leading `s3:` so the worker accepts
MinIO events without a MinIO-shaped escape hatch in the AWS-S3 normalizer.
3. The python summary in run-tests.sh collapsed every non-PASS status to
`[FAIL]`, hiding the distinction between actual assertion failures and
probe-skipped ERRORs. Render the real status so triage isn't misleading.
Verified end-to-end: bash run-tests.sh --providers=local,s3 now reports
pass=64 fail=0 error=0; all 5 s3 trigger scenarios pass (object-created
fires on putObject, object-deleted fires on deleteObject, multi-put
coalescing, getObject silence, rapid-fire 10 puts). bash run-tests.sh
--providers=all reports pass=89 fail=0 error=5 — the 5 ERRORs are the r2
trigger cases skipping cleanly because R2 trigger plumbing (Cloudflare
Queue) is intentionally not wired in the harness.
When --providers=all is used (CI default), r2 trigger cases ERROR-skip because there's no Cloudflare Queue plumbing in this harness — the sentinel embeds "ERROR" and run-tests.sh exits 2, which fails the CI job for an intentional gap rather than a real regression. Add a HARNESS_TRIGGER_PROVIDERS env (and matching --trigger-providers= flag) that constrains which providers get trigger registration, probe, and case fan-out. Defaults to all selected providers, so existing callers and developer flows are unchanged. CI workflow now passes --trigger-providers=local,s3 alongside --providers=all. R2 RPC, edge, quirk, and concurrency suites still run (scratch-r2 against MinIO — 25 cases preserved), but the 5 r2 trigger scenarios are no longer fanned out, so the run exits 0 cleanly. Verified: bash run-tests.sh --providers=all --trigger-providers=local,s3 reports pass=89 fail=0 error=0; exit 0. by_provider.r2 = 25/0/0 (RPC suite intact); by_provider.local = 29/0/0 and by_provider.s3 = 30/0/0 (trigger fan-out covers both).
Closes the gaps surfaced by the compliance review against
workers/binary-worker.md.
Publish-blocking (CRITICAL):
- Add build.rs exposing the build-time TARGET triple.
- Add src/manifest.rs with build_manifest() emitting the five fields
POST /publish requires (name, version, description, default_config,
supported_targets) plus JSON-roundtrip unit tests.
- Add --manifest CLI flag with early-exit branch in main.rs so the
registry publish pipeline can extract the manifest without an engine
connection.
- Add tests/manifest.rs (spec section 9 pattern A) that spawns the
binary with --manifest and validates the JSON contract.
Structural (HIGH):
- Pass WorkerMetadata { runtime, version, name, os, pid, telemetry }
to register_worker so operators and the registry get a stable
identity line.
- Commit a default config.yaml that wires a rustfs-backed scratch
bucket so a fresh checkout runs end-to-end with no credentials.
- Replace the fatal config-load with a warn+default fallback (config
errors no longer crash the worker).
- Move inline register_function calls into
handlers::register_all(iii, &state); main.rs now has one
registration call instead of forty lines.
- Split tests/integration.rs: the schema-regression test moves to
tests/schemas.rs, and tests/integration.rs becomes the spec
pattern-A subprocess harness with self-skip when iii or rustfs is
missing and a readiness retry around the first putObject.
Polish (MEDIUM/LOW):
- Add free pub fn load_config(path) -> anyhow::Result<WorkerConfig>
matching the spec template signature.
- Derive Default on WorkerConfig (enables the graceful fallback above).
- Migrate serde_yml -> serde_yaml = "0.9" for spec parity.
- Add "Local development & testing" section to README with run/test
commands and the section 11 preflight checklist.
Verified from workers/storage/:
cargo fmt --all -- --check OK
cargo clippy --all-targets --all-features -- -D warnings OK
cargo test --lib --test manifest --test schemas --test integration
-> 81 passed, 3 ignored
./target/debug/storage --manifest | jq . prints all 5 fields
Drive-by reformatting from running `cargo fmt --all` as part of the spec-compliance preflight in the previous commit. No logic changes: just line breaks and parameter wrapping that rustfmt prefers.
025d976 to
f89d527
Compare
There was a problem hiding this comment.
Actionable comments posted: 15
Note
Due to the large number of review comments, Critical, Major severity comments were prioritized as inline comments.
🟡 Minor comments (13)
storage/README.md-54-71 (1)
54-71:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winRust SDK example incorrectly uses
TriggerRequestto invoke an RPC.
iii.trigger(TriggerRequest { function_id: "storage::putObject" ... })dispatches a trigger event, not an RPC call. Invokingstorage::putObject(an RPC handler) via the trigger API will not behave as shown. The Rust SDK equivalent of the TypeScriptcall('storage::putObject', ...)should use a call/invoke method, nottrigger.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@storage/README.md` around lines 54 - 71, The example incorrectly uses iii.trigger with a TriggerRequest to call the RPC handler storage::putObject; change the example to use the Rust SDK's RPC/call/invoke method (the counterpart to TypeScript's call('storage::putObject', ...)) instead of TriggerRequest/iii.trigger so the RPC is invoked properly; locate usage of TriggerRequest and iii.trigger in the README example and replace it with the SDK's call/invoke API, passing the function id "storage::putObject" and the same payload and timeout parameters.storage/src/rustfs/health.rs-20-27 (1)
20-27:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winDeadline can be exceeded by a full request timeout.
The loop can run a final
.send()with a 2s timeout even when almost no budget remains, somax_waitis not strictly honored. Clamp each attempt timeout to remaining time.Suggested fix
pub async fn wait_for_healthy(port: u16, max_wait: Duration) -> Result<(), StorageError> { let deadline = Instant::now() + max_wait; let url = format!("http://127.0.0.1:{port}/"); - let client = reqwest::Client::builder() - .timeout(Duration::from_secs(2)) + let client = reqwest::Client::builder() .build() .map_err(|e| StorageError::LocalBackendBootFailed { reason: format!("reqwest build: {e}"), })?; while Instant::now() < deadline { - if let Ok(resp) = client.get(&url).send().await { + let remaining = deadline.saturating_duration_since(Instant::now()); + let req_timeout = remaining.min(Duration::from_secs(2)); + if req_timeout.is_zero() { + break; + } + if let Ok(resp) = client.get(&url).timeout(req_timeout).send().await { if resp.status().as_u16() < 500 { return Ok(()); } } sleep(Duration::from_millis(200)).await;🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@storage/src/rustfs/health.rs` around lines 20 - 27, The loop can exceed max_wait because each iteration calls client.get(&url).send().await with a fixed request timeout; clamp each attempt to the remaining time by computing let remaining = deadline.saturating_duration_since(Instant::now()) and using either client.get(&url).timeout(remaining.min(MAX_REQ_TIMEOUT)).send().await or wrapping the send() in tokio::time::timeout(remaining.min(MAX_REQ_TIMEOUT), client.get(&url).send()).await, then handle the timeout/error case the same as other errors and continue/sleep; reference the existing symbols deadline, Instant::now(), client.get(&url).send().await, and sleep(Duration::from_millis(...)).await.storage/tests/e2e/workers/harness/src/cases-concurrency.ts-17-19 (1)
17-19:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winStrengthen ETag assertions to catch missing response fields.
etags.size >= 1is too weak here and can still pass when all puts return missing/undefinedetag.Suggested patch
- const etags = new Set(puts.map((p) => p.etag as string)); - assertTruthy(etags.size >= 1, `expected at least one etag, got: ${JSON.stringify([...etags])}`); + const etags = new Set( + puts.map((p) => p.etag).filter((e): e is string => typeof e === 'string' && e.length > 0), + ); + assertTruthy(etags.size === puts.length, + `expected etag on every put, got ${etags.size}/${puts.length}`);🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@storage/tests/e2e/workers/harness/src/cases-concurrency.ts` around lines 17 - 19, The current assertion using etags.size >= 1 is too weak and can pass when many put responses are missing etag; update the check around the etags Set created from puts.map(...) and the assertTruthy call to verify every put has a defined, non-empty etag (e.g., assert that puts.every(p => p.etag != null && p.etag !== '') and that etags.size === puts.length) so missing/undefined etag fields fail the test; refer to the etags variable, the puts array mapping, and the assertTruthy invocation to locate and update the assertion.storage/tests/e2e/workers/harness/src/worker.ts-17-27 (1)
17-27:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winValidate provider names instead of blind casting.
as Provider[]trusts env input. Unknown provider strings should fail fast with a clear message at parse time.Suggested patch
const PROVIDERS_RAW = process.env.HARNESS_PROVIDERS ?? 'local'; -const PROVIDERS = PROVIDERS_RAW.split(',').map((s) => s.trim()).filter(Boolean) as Provider[]; +const ALLOWED_PROVIDERS = new Set<Provider>(['local', 's3', 'gcs', 'r2']); +const PROVIDERS = PROVIDERS_RAW.split(',') + .map((s) => s.trim()) + .filter(Boolean) + .map((p) => { + if (!ALLOWED_PROVIDERS.has(p as Provider)) { + throw new Error(`Invalid provider in HARNESS_PROVIDERS: ${p}`); + } + return p as Provider; + }); @@ const TRIGGER_PROVIDERS: Provider[] | undefined = TRIGGER_PROVIDERS_RAW - ? (TRIGGER_PROVIDERS_RAW.split(',').map((s) => s.trim()).filter(Boolean) as Provider[]) + ? TRIGGER_PROVIDERS_RAW.split(',') + .map((s) => s.trim()) + .filter(Boolean) + .map((p) => { + if (!ALLOWED_PROVIDERS.has(p as Provider)) { + throw new Error(`Invalid provider in HARNESS_TRIGGER_PROVIDERS: ${p}`); + } + return p as Provider; + }) : undefined;🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@storage/tests/e2e/workers/harness/src/worker.ts` around lines 17 - 27, The code blindly casts environment strings to Provider[] via PROVIDERS_RAW/PROVIDERS and TRIGGER_PROVIDERS_RAW/TRIGGER_PROVIDERS which can hide invalid provider names; update parsing to validate each token against the canonical set of allowed Provider values (e.g., an enum or array of valid provider strings) rather than using "as Provider[]", and if any token is not recognized throw a clear, early error (with the invalid names listed) during initialization; factor this into a small helper like parseProviders(raw: string): Provider[] and use it for both PROVIDERS and TRIGGER_PROVIDERS to ensure fail-fast validation.storage/tests/e2e/workers/harness/src/cases-provider.ts-17-35 (1)
17-35:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winActually validate URL parseability in the S3 presign case.
The case name/intent says “URL is parseable,” but current assertion only checks
includes(bucket).Suggested patch
const r = await ctx.call('storage::presignUrl', { bucket: ctx.bucket, key, method: 'GET', expires_in_seconds: 60, }); + let parsed: URL; + try { + parsed = new URL(r.url); + } catch { + throw new Error(`presigned URL is not parseable: ${r.url}`); + } @@ - assertTruthy(r.url.includes(ctx.bucket), - `presigned URL should reference bucket ${ctx.bucket}: ${r.url}`); + assertTruthy( + parsed.hostname.includes(ctx.bucket) || parsed.pathname.includes(`/${ctx.bucket}/`), + `presigned URL should reference bucket ${ctx.bucket}: ${r.url}`, + );🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@storage/tests/e2e/workers/harness/src/cases-provider.ts` around lines 17 - 35, The test's presign URL assertion only checks r.url.includes(ctx.bucket) but must actually validate URL parseability and that the bucket appears either in hostname or pathname; update the run(ctx) for the '[s3] presigned URL is path-style against MinIO' case to attempt parsing r.url with the URL constructor (ensuring it doesn't throw), then assert that either parsed.hostname or parsed.pathname contains ctx.bucket (so both virtual-host and path-style URLs pass), using the existing r from the storage::presignUrl call to produce clear assertion messages if parsing fails or the bucket is not found.storage/tests/e2e/s3.rs-20-24 (1)
20-24:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winSkip early when AWS credentials are unset.
This test currently skips only on
S3_TEST_BUCKET; missingAWS_ACCESS_KEY_ID/AWS_SECRET_ACCESS_KEYwill fail later with less actionable auth errors.Suggested patch
async fn s3_round_trip() { let Some(bucket) = require_env("S3_TEST_BUCKET") else { return; }; + let Some(_) = require_env("AWS_ACCESS_KEY_ID") else { + return; + }; + let Some(_) = require_env("AWS_SECRET_ACCESS_KEY") else { + return; + }; let region = std::env::var("S3_TEST_REGION").unwrap_or_else(|_| "us-east-1".into());🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@storage/tests/e2e/s3.rs` around lines 20 - 24, The test currently only early-returns when S3_TEST_BUCKET is missing, causing unclear auth failures later; update the startup checks to also require AWS_ACCESS_KEY_ID and AWS_SECRET_ACCESS_KEY (e.g., call require_env for "AWS_ACCESS_KEY_ID" and "AWS_SECRET_ACCESS_KEY" along with "S3_TEST_BUCKET" before constructing BucketConfig/S3BucketConfig) and return early if either credential is unset so the test is skipped with a clear reason.storage/tests/e2e/fixtures/minio-sqs-bridge/index.ts-77-81 (1)
77-81:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winGuard against non-object records before dereferencing
Line 78-80 assumes each
Recordselement is object-shaped;nullor primitive values can throw and return 500. Filter/guard invalid elements before readingeventName.Suggested fix
- const normalized = body.Records.map((rec) => { - const r = rec as Record<string, unknown> + const normalized = body.Records + .filter((rec): rec is Record<string, unknown> => !!rec && typeof rec === 'object') + .map((r) => { const en = typeof r.eventName === 'string' ? r.eventName : '' return en.startsWith('s3:') ? { ...r, eventName: en.slice('s3:'.length) } : r - }) + })🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@storage/tests/e2e/fixtures/minio-sqs-bridge/index.ts` around lines 77 - 81, The code mapping body.Records into normalized assumes each record is a non-null object and directly dereferences rec.eventName; update the logic around normalized to first filter out non-object or null records (e.g., ensure typeof rec === 'object' && rec !== null) and only then cast to Record<string, unknown> before reading eventName, and handle missing or non-string eventName safely so primitives/nulls are skipped or returned unchanged; change references in this block (body.Records and normalized) accordingly.storage/tests/e2e/workers/harness/src/runner.ts-181-187 (1)
181-187:⚠️ Potential issue | 🟡 Minor | ⚡ Quick win
resetEventsclears waiter timers without rejecting them — any in-flightwaitForEventwill hang forever.If a case (or future case) calls
ctx.resetEvents()while awaitForEventis still pending — the most likely path is a case that doesPromise.all([waitForEvent(...), somethingElse()])and then resets in afinally— the promise loses both its resolver (we drop the entry) and its timeout-rejector (weclearTimeoutit), leaving the awaiter wedged indefinitely. One-line fix: reject before dropping.🛡️ Reject pending waiters on reset
private resetEvents(): void { - // Cancel any pending waiters first; their cases either already failed - // or are no longer interested. - for (const w of this.waiters) clearTimeout(w.timer); + // Cancel any pending waiters first; their cases either already failed + // or are no longer interested. Reject so awaiters don't leak. + for (const w of this.waiters) { + clearTimeout(w.timer); + w.reject(new Error('waitForEvent aborted by resetEvents')); + } this.waiters = []; this.events = []; }🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@storage/tests/e2e/workers/harness/src/runner.ts` around lines 181 - 187, resetEvents currently clears waiter timers and drops waiter entries causing any in-flight waitForEvent promises to hang; modify resetEvents (which manipulates this.waiters, this.events and each waiter.timer) to first reject each pending waiter (call its rejector with a clear Error/Cancel reason), then clear its timer, and finally clear this.waiters and this.events so no promise is left unresolved; ensure you reference the waiter objects used by waitForEvent so their resolve/reject callbacks are invoked before removal.storage/src/backend/gcs.rs-200-219 (1)
200-219:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winPresign URLs ignore
endpoint_urland always targetstorage.googleapis.com.The build-time comment (lines 60–64) acknowledges this, but the
presignmethod itself silently returns a real-GCS URL when the bucket is configured against a custom endpoint (e.g.,fake-gcs-serverfor the e2e harness, or any private GCS-compatible store). Callers that hand the URL back to a client will see "host not reachable" / signature mismatch errors rather than a clear configuration message.Consider rejecting
presignupfront when an override is configured, with a documented error so harness/operators get a deterministic failure rather than a runtime 4xx from the wrong host:🛠 Suggested guard at presign entry
// At GcsBackend::build, stash endpoint_url.is_some() into a field, then: async fn presign(&self, req: PresignReq) -> Result<PresignResp, BackendError> { if self.endpoint_overridden { return Err(BackendError::PresignUnsupported( "gcs presign always signs against storage.googleapis.com; \ not supported when endpoint_url is overridden (fake-gcs-server, etc.)".into() )); } // ... }🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@storage/src/backend/gcs.rs` around lines 200 - 219, The presign method currently always signs against storage.googleapis.com and must early-fail when a custom endpoint is configured; add a boolean field (e.g., endpoint_overridden) to the GcsBackend struct during GcsBackend::build to record endpoint_url.is_some(), then modify async fn presign(&self, req: PresignReq) to check that flag and return Err(BackendError::PresignUnsupported("gcs presign always signs against storage.googleapis.com; not supported when endpoint_url is overridden (fake-gcs-server, etc.)".into())) when true, otherwise proceed with the existing signing logic.storage/src/triggers/pollers/sqs.rs-94-104 (1)
94-104:⚠️ Potential issue | 🟡 Minor | ⚡ Quick win
delete_messageerrors silently dropped.
let _ = self.client.delete_message()...send().await;discards both Ok and Err. If the delete fails (network blip, expired receipt handle, permissions), the message redelivers and you get duplicate dispatches with no signal in logs that anything went wrong. Awarn!here would make that visible without changing semantics.🛠 Proposed fix
- let _ = self - .client - .delete_message() - .queue_url(&self.queue_url) - .receipt_handle(rh) - .send() - .await; + if let Err(e) = self + .client + .delete_message() + .queue_url(&self.queue_url) + .receipt_handle(rh) + .send() + .await + { + tracing::warn!(error=%e, "sqs delete_message failed; message will redeliver"); + }🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@storage/src/triggers/pollers/sqs.rs` around lines 94 - 104, The delete_message() call currently discards the Result from send().await; change it to inspect the result and emit a warn! when it fails (include the error, the receipt_handle value and self.queue_url for context) while keeping the same non-panicking behavior; locate the block handling msg.receipt_handle inside the poller (where delete_message(), .queue_url(&self.queue_url) and .receipt_handle(rh) are invoked) and replace the `let _ = ...send().await` with a match or if let Err(e) => warn!(...) that logs e and identifying info.storage/tests/e2e/run-tests.sh-335-353 (1)
335-353:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winHealth-wait loop misclassifies services without a
healthcheckas unhealthy.
docker compose ps --format '{{.Service}} {{.Health}}'emits an empty{{.Health}}field for services that don't declare ahealthcheck:block. The awk filter$2 != "healthy"then matches those rows and treats them as unhealthy forever, which can hang the loop until the 60s deadline and then produce a misleading FATAL listing services that never had a healthcheck to begin with.Either require every compose service to define a healthcheck, or treat empty as "ok":
🛠 Proposed fix
unhealthy=$(cd "$ROOT_DIR" && \ docker compose --profile cloud ps --format '{{.Service}} {{.Health}}' \ - | awk '$2 != "healthy" {print $1}') + | awk '$2 != "" && $2 != "healthy" {print $1}')🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@storage/tests/e2e/run-tests.sh` around lines 335 - 353, The health-wait loop misclassifies services with no healthcheck as unhealthy; update the awk filter in the unhealthy assignment (the line that runs docker compose --profile cloud ps --format '{{.Service}} {{.Health}}' and pipes to awk) so it ignores empty {{.Health}} fields (treat empty as OK) — e.g. only print services whose second field is present and not "healthy" — leaving the rest of the loop (health_deadline, break/exit logic, and per-service logs) unchanged.storage/src/triggers/pollers/pubsub.rs-99-106 (1)
99-106:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winSilent JSON-parse fallback can produce a permanent nack loop on malformed messages.
When
pm.dataisn't valid JSON,databecomesNull.normalize_gcswill then return a non-Ignorederror, which the dispatcher path translates to "do not ack" → Pub/Sub redelivers → same parse failure → forever. For poison messages we want to drop (ack), not retry.Suggest logging at warn and acking on parse failure (or returning
NormalizeError::Ignoredupstream so the existing ack-on-Ignored branch handles it):🛠 Proposed fix
- let data = serde_json::from_slice::<serde_json::Value>(&pm.data) - .unwrap_or(serde_json::Value::Null); + let data = match serde_json::from_slice::<serde_json::Value>(&pm.data) { + Ok(v) => v, + Err(e) => { + tracing::warn!( + message_id = %pm.message_id, + error = %e, + "pubsub message body not JSON; acking to drop" + ); + if let Err(e) = received.ack().await { + tracing::warn!(error = %e, "pubsub ack of poison message failed"); + } + continue; + } + };🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@storage/src/triggers/pollers/pubsub.rs` around lines 99 - 106, The current silent JSON-parse fallback turns malformed pm.data into Null and causes normalize_gcs to emit a non-Ignored error that leads to endless redelivery; replace the unwrap_or handling so a JSON parse Err triggers a warning log (include pm.message_id and error info) and returns NormalizeError::Ignored (so the dispatcher will ack), e.g. change the serde_json::from_slice::<Value>(&pm.data).unwrap_or(...) usage in pubsub.rs to a match/if let that logs warn and returns early with NormalizeError::Ignored before constructing the envelope; this ensures normalize_gcs/dispatcher ack poisoned messages instead of retrying forever.storage/src/rustfs/spawn.rs-174-181 (1)
174-181:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winSerialization recommended for env-var mutations in tests.
std::env::set_var/remove_varmutate process-global state. Sincediscover_binary()readsRUSTFS_BINand tests run in parallel by default, this test can race with concurrent execution. Recommend serializing with theserial_testcrate or usingtemp_env::with_var(...)for scoped, thread-safe overrides.Note:
std::env::set_varandremove_varwill becomeunsafein Rust 2024 edition due to inherent unsoundness in multithreaded contexts; migration totemp_envsidesteps this future concern.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@storage/src/rustfs/spawn.rs` around lines 174 - 181, The test function discover_binary_reads_env_var mutates global process env with std::env::set_var/remove_var which races when tests run in parallel; update the test (discover_binary_reads_env_var) to scope the environment change using a thread-safe helper such as temp_env::with_var or annotate/serialize the test using the serial_test crate (e.g., #[serial]) so RUSTFS_BIN is set only for the test duration and automatically restored/isolated, ensuring discover_binary() reads the temporary value without global races.
🧹 Nitpick comments (8)
storage/build.rs (1)
4-4: 💤 Low valuePrefer
expectoverunwrapfor a more actionable build failure message.
TARGETis always set by Cargo in build scripts, so this won't panic in practice, butexpectcommunicates intent more clearly.🔧 Suggested improvement
- std::env::var("TARGET").unwrap() + std::env::var("TARGET").expect("Cargo always sets TARGET in build scripts")🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@storage/build.rs` at line 4, Replace the unwrap() call on std::env::var("TARGET") in build.rs with expect() to provide a clearer, actionable error message; update the expression where std::env::var("TARGET").unwrap() appears so it uses expect("TARGET must be set by Cargo build scripts") (or similar) to document intent and improve the build-time panic message..github/workflows/storage-e2e.yml (1)
64-68: ⚡ Quick winPin the
iiiengine version to prevent non-reproducible CI runs.Installing from
mainmeans every CI run may pick up a differentiiiengine version. A breaking change in the engine will silently flip all e2e tests from green to red with no code change in this repo.🔧 Suggested fix
+env: + CARGO_TERM_COLOR: always + III_VERSION: '0.x.y' # pin to a known-good release tag + ... - name: Install iii engine run: | - curl -fsSL --retry 3 --retry-connrefused --retry-delay 5 \ - https://install.iii.dev/iii/main/install.sh | sh + curl -fsSL --retry 3 --retry-connrefused --retry-delay 5 \ + "https://install.iii.dev/iii/${III_VERSION}/install.sh" | sh echo "$HOME/.local/bin" >> "$GITHUB_PATH"🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In @.github/workflows/storage-e2e.yml around lines 64 - 68, The workflow step "Install iii engine" currently pulls install.sh from the main branch which causes non-reproducible CI; update that step to pin the installer to a specific, committed release/tag or version (e.g., replace the "main" URL with a release/tag or add an explicit version parameter/ENV used by the installer) so every run installs the same iii engine version; ensure the step name "Install iii engine" and the curl invocation are updated accordingly and include a comment or variable (e.g., III_VERSION) so future updates require an explicit version bump.storage/src/triggers/dispatcher.rs (1)
36-111: ⚖️ Poor tradeoffLGTM — ack semantics and timeout handling line up with the iii-database contract.
Worth calling out for a future revision: the per-subscriber
awaitis sequential, so total dispatch latency for a(bucket, kind)with N subscribers isΣ(handler_timeout_ms+1s). Pollers (SQS visibility timeout, CF Queue lease) only get one ack decision for the message, so a single slow subscriber stalls the entire group. With the current 1-subscriber-typical shape this is fine, but if/when fan-out grows, switching tofutures::future::join_allover the per-sub futures would bound the wall-clock atmax(...)instead ofsum(...).🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@storage/src/triggers/dispatcher.rs` around lines 36 - 111, Sequential per-subscriber awaits in dispatch cause total latency to sum per-subscriber time; change dispatch in storage::triggers::dispatcher.rs to fire all per-subscriber trigger futures concurrently (e.g., map subs from registry.subscribers_for(...) to individual futures that call self.iii.trigger(req) wrapped with the same timeout logic using sub.handler_timeout_ms) and then await them together with futures::future::join_all or similar, preserving the existing ack extraction logic (null => true, value.get("ack")... ) and keeping function_id and timeout_ms logging behavior.storage/src/triggers/pollers/rustfs_webhook.rs (1)
32-68: ⚡ Quick win
spawn_receiverandspawn_receiver_onare line-for-line duplicates apart from the bind argument.Both functions build the same router, channel, and
axum::servetask; only theTcpListener::bindcall differs. Easier to maintain (and harder to drift) ifspawn_receiverdelegates throughspawn_receiver_on(state, 0)—127.0.0.1:0already gives the ephemeral-port behavior the doc-comment promises.♻️ Collapse to one bind helper
-pub async fn spawn_receiver(state: WebhookState) -> std::io::Result<WebhookHandle> { - let listener = TcpListener::bind("127.0.0.1:0").await?; - let addr = listener.local_addr()?; - let app = Router::new() - .route("/notify", post(notify_handler)) - .with_state(state); - let (shutdown_tx, shutdown_rx) = oneshot::channel(); - tokio::spawn(async move { - let _ = axum::serve(listener, app) - .with_graceful_shutdown(async { - let _ = shutdown_rx.await; - }) - .await; - }); - Ok(WebhookHandle { addr, shutdown_tx }) -} - -/// Same as `spawn_receiver` but binds on a pre-allocated port. ... -pub async fn spawn_receiver_on(state: WebhookState, port: u16) -> std::io::Result<WebhookHandle> { - let listener = - TcpListener::bind(SocketAddr::from((std::net::Ipv4Addr::LOCALHOST, port))).await?; - let addr = listener.local_addr()?; - let app = Router::new() - .route("/notify", post(notify_handler)) - .with_state(state); - let (shutdown_tx, shutdown_rx) = oneshot::channel(); - tokio::spawn(async move { - let _ = axum::serve(listener, app) - .with_graceful_shutdown(async { - let _ = shutdown_rx.await; - }) - .await; - }); - Ok(WebhookHandle { addr, shutdown_tx }) -} +pub async fn spawn_receiver(state: WebhookState) -> std::io::Result<WebhookHandle> { + spawn_receiver_on(state, 0).await +} + +/// Bind the loopback HTTP receiver on `port` (use `0` for an ephemeral port). +pub async fn spawn_receiver_on(state: WebhookState, port: u16) -> std::io::Result<WebhookHandle> { + let listener = + TcpListener::bind(SocketAddr::from((std::net::Ipv4Addr::LOCALHOST, port))).await?; + let addr = listener.local_addr()?; + let app = Router::new() + .route("/notify", post(notify_handler)) + .with_state(state); + let (shutdown_tx, shutdown_rx) = oneshot::channel(); + tokio::spawn(async move { + let _ = axum::serve(listener, app) + .with_graceful_shutdown(async { + let _ = shutdown_rx.await; + }) + .await; + }); + Ok(WebhookHandle { addr, shutdown_tx }) +}🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@storage/src/triggers/pollers/rustfs_webhook.rs` around lines 32 - 68, spawn_receiver duplicates spawn_receiver_on except for the bind argument; change spawn_receiver to simply call and return spawn_receiver_on(state, 0) so the ephemeral-port behavior is delegated to spawn_receiver_on, eliminating the duplicated router/serve/channels code (keep the same async signature and return type for spawn_receiver and remove the duplicate implementation).storage/src/triggers/handler.rs (1)
15-98: 💤 Low valueTwo near-identical handlers — consider a single generic over
EventKind.
ObjectCreatedHandlerandObjectDeletedHandlerdiffer only in (a) theTriggerConfignewtype they parse, (b) theEventKindthey register, and (c) the log string. A single generic struct (or shared private helper that takesEventKind+ a parser closure) would halve the surface area and ensure the two paths can't drift (e.g., the JSON-injection fix above has to be applied twice). Optional given chill-mode and that the file is small.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@storage/src/triggers/handler.rs` around lines 15 - 98, The two handlers ObjectCreatedHandler and ObjectDeletedHandler duplicate register_trigger/unregister_trigger logic; consolidate into one generic handler to avoid drift. Replace both with a single struct (e.g., ObjectEventHandler) parameterized by EventKind and a parser function/enum specifying which config type to deserialize (CreatedConfig vs DeletedConfig) so register_trigger calls the chosen parser, validates wired_buckets, calls registry.register with the injected EventKind, and emits a log message derived from that EventKind; keep unregister_trigger delegating to registry.unregister(&config.id) unchanged. Implement the parser injection as a Fn(&Value) -> Result<ParsedConfig, IIIError> or an enum branch inside the generic handler to locate logic around register_trigger in ObjectCreatedHandler/ObjectDeletedHandler and switch to the unified register_trigger implementation.storage/tests/e2e/workers/harness/src/cases.ts (2)
163-168: 💤 Low valueConfusing assertion failure message.
assertEqual(del.deleted, true, 'delete returned deleted=false')reads as if the failure mode isdeleted=false, butassertEqualonly fires when the value differs fromtrue. On failure the user sees the string "delete returned deleted=false" alongside the actual mismatch — fine in this specific case but inverted phrasing for anassertEqual(..., true, ...). Recommend'expected delete.deleted=true'.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@storage/tests/e2e/workers/harness/src/cases.ts` around lines 163 - 168, The assertion message for the delete check is misleading: update the assertion in the block that calls ctx.call('storage::deleteObject') so the assertEqual on del.deleted uses a positive expectation message (e.g., "expected del.deleted=true" or "expected delete.deleted to be true") instead of "delete returned deleted=false"; target the assertEqual call that compares del.deleted to true to make the message reflect the expected state.
41-41: 💤 Low valueTypo in JSDoc: "an storage" → "a storage".
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@storage/tests/e2e/workers/harness/src/cases.ts` at line 41, The JSDoc above the call property contains a typo "an storage" — update the comment to read "a storage" (or reword to "the storage" if more appropriate) for the call: (functionId: string, payload: unknown) => Promise<any>; identifier so the documentation is grammatically correct.storage/tests/e2e/run-tests.sh (1)
253-272: 💤 Low valueSHA256 verification is best-effort even when checksums file is reachable but unparseable.
If
SHA256SUMSis reachable butawkfinds no line for$asset, the script logsWARN ... skipping verificationand continues. For a CI flow that auto-fetches a binary over the public internet, silent fallback widens the supply-chain surface — a cache poisoning that returns a differentSHA256SUMSwould not block download. Consider failing closed (exit 1) when the asset is expected to be listed but isn't, gated behind aRUSTFS_REQUIRE_SHA256=1env so local dev keeps the soft path.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@storage/tests/e2e/run-tests.sh` around lines 253 - 272, The current SHA256 verification silently skips when the SHA256SUMS file is reachable but awk finds no entry for "$asset"; change this to fail when the caller requires strict verification by checking an env var (RUSTFS_REQUIRE_SHA256=1): in the run-tests.sh block that computes expected="$(awk -v a="$asset" '$2==a {print $1}' "$tmpdir/SHA256SUMS")", if [[ -z "$expected" ]] then if RUSTFS_REQUIRE_SHA256 is set to 1 exit with an error (return 1) and log a clear fatal message referencing "$asset" and "$tmpdir/SHA256SUMS", otherwise keep the existing warn-and-skip behavior; preserve the existing sha256sum/shasum branch and subsequent checks for actual vs expected.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@storage/src/backend/factory.rs`:
- Around line 38-48: S3 path-style routing is hardcoded false in factory.rs when
calling S3Backend::build via S3BackendOpts, causing S3-compatible endpoints to
fail; add a force_path_style: bool field to S3BucketConfig
(storage/src/config.rs), read it where configs are loaded, and pass that value
into S3BackendOpts.force_path_style in factory.rs instead of the literal false
(update any call sites such as the one inside the S3Backend::build invocation
and ensure tests like build_with_endpoint_url_does_not_panic use the new config
field).
In `@storage/src/backend/local.rs`:
- Around line 46-58: The current logic treats any head_bucket() failure as
"bucket missing" and always calls create_bucket(); modify the flow so
head_bucket().send().await is matched: if Ok(_) return Ok(()); if Err(e) inspect
the error and only proceed to client.create_bucket() when the error indicates a
definite NotFound/NoSuchBucket (HTTP 404 or the provider-specific NotFound
variant); for any other error (auth, connection, etc.) immediately map and
return a BackendError::Provider with the original error info (preserve e in the
message), and keep the create_bucket() call and its error mapping unchanged so
creation only happens on explicit not-found errors.
- Line 92: The call to std::fs::create_dir_all(queue_dir) currently ignores
errors and should fail fast; update the code that constructs the queue_dir (the
create_dir_all(queue_dir) call) to handle the Result instead of discarding
it—propagate the error (use the function's error return path) or return an Err
with context so registration stops when the directory cannot be created; ensure
you use the existing error type or map the std::io::Error into the function's
error type so callers (and tests) see a clear failure instead of silently
continuing.
In `@storage/src/backend/s3.rs`:
- Around line 121-141: The code currently treats
resp.content_length.unwrap_or(0) as 0 which lets objects with unknown length
bypass req.max_inline_bytes and then be fully buffered by resp.body.collect();
change the logic so that before calling resp.body.collect() you first check if
req.max_inline_bytes.is_some() and if resp.content_length.is_none() fail closed
(return a BackendError::ObjectTooLarge or a clear Provider error) OR
alternatively enforce the cap while streaming (read resp.body incrementally and
abort with BackendError::ObjectTooLarge when the cumulative bytes exceed cap);
update the branch around resp.content_length, the req.max_inline_bytes check,
and the resp.body.collect() usage to implement one of these two approaches.
In `@storage/src/handlers/put_object.rs`:
- Line 16: The PutObject request struct only accepts body_base64 which breaks
compatibility with clients sending body; update the PutObject request field (the
struct containing pub body_base64: String, e.g., PutObjectRequest) to accept the
legacy key by adding a serde alias for "body" on that field (use #[serde(alias =
"body")] or equivalent deserialization attribute) so incoming payloads with
"body" or "body_base64" both deserialize correctly; run/update relevant tests to
cover both keys.
- Around line 37-47: Reject oversized payloads before allocating by checking the
base64 string length instead of decoding first: in put_object.rs, before calling
B64.decode(&req.body_base64) compute the maximum decoded byte length from
req.body_base64 (e.g. (req.body_base64.len() * 3) / 4 minus padding) and compare
against INLINE_BODY_CAP, returning the StorageError::BodyTooLarge via err_to_str
if it exceeds cap; only call B64.decode(&req.body_base64) afterwards and keep
existing InvalidBase64 handling for decode errors.
In `@storage/src/triggers/handler.rs`:
- Around line 26-36: The error JSON is being constructed by interpolating
Display values (serde_json::Error and cfg.bucket) into a raw string which can
produce malformed JSON; instead build the payload as a serde_json value and
serialize it so values are properly escaped — replace the raw r#"...{e}..."# and
r#"...{}`..."# constructions used when returning IIIError::Handler in the
object-created handler (where CreatedConfig is parsed and cfg.bucket is checked)
with a serde_json::json! or serde_json::Map that sets
{"code":"CONFIG_ERROR","message": <string>} and call to_string() for the error
envelope; apply the identical fix to the ObjectDeletedHandler::register_trigger
block (the duplicate at lines ~67–77) so both handlers produce safe JSON-escaped
messages.
In `@storage/src/triggers/normalize.rs`:
- Around line 18-19: The decode_s3_key function currently only percent-decodes
%XX sequences and misses S3-style '+' space encoding; update decode_s3_key to
first replace all '+' characters with space characters (or decode '+' to ' ')
before calling percent_decode_str so inputs like "reports+2026.pdf" normalize
correctly and match RPC-sourced keys; reference the decode_s3_key function and
ensure the replacement occurs prior to percent_decode_str().
In `@storage/src/triggers/pollers/cf_queue.rs`:
- Around line 191-197: probe_auth currently only errors on 401/403 but returns
Ok(()) for 404/429/5xx; change it to treat any non-2xx as a startup failure by
checking resp.status().is_success() and returning an Err when false. Keep the
existing special-case for reqwest::StatusCode::UNAUTHORIZED / FORBIDDEN
(StorageError::CfQueueAuthFailed with the current message), and for all other
non-success statuses return a suitable StorageError (e.g., a generic
CfQueueAuthFailed or a new variant) that includes the HTTP status and response
body or context so startup fails when the probe returns 404/429/5xx. Ensure you
update probe_auth to use resp.status().is_success(), reference resp and status,
and include the status details in the error message.
- Around line 115-124: The match on normalize_r2 currently drops non-ignored
errors without acking, causing permanent malformed (poison) messages to be
retried forever; update the match in the poller handling (the normalize_r2(...)
match block) so that for non-ignored/unmapped and for any other Err(_) cases you
log the error with context (payload and NormalizeError) and push lease_id into
acked_lease_ids (or otherwise dead-letter it) before continue; ensure you still
preserve the existing behavior for NormalizeError::Ignored and
NormalizeError::UnmappedBucket where acking is already done, and only add
ack+log for permanently malformed cases so the lease is removed instead of
retried.
In `@storage/src/triggers/pollers/sqs.rs`:
- Around line 71-81: The normalize_s3 call currently continues the SQS message
loop on NormalizeError::Ignored(_) and NormalizeError::UnmappedBucket(_),
leaving the SQS message un-acked and causing infinite redelivery; change those
branches to ack-and-drop the SQS message instead of continuing: mirror
pubsub.rs::dispatch_one behavior by performing the same acknowledgement (invoke
the SQS delete/ack path or return the same "ack" sentinel used by this poller)
when encountering NormalizeError::Ignored and NormalizeError::UnmappedBucket,
and keep the existing tracing/logging; leave the other Err(e) branch as-is.
In `@storage/tests/e2e/fixtures/minio-sqs-bridge/Dockerfile`:
- Around line 3-23: The Dockerfile runs the Bun service as root; add an explicit
non-root runtime user and switch to it before CMD to harden the container: in
the Dockerfile create a dedicated user/group (e.g., app or bunuser) after
setting WORKDIR, ensure ownership of /app and any copied files is changed to
that user (chown), and then set USER to that non-root account so HEALTHCHECK and
CMD ["bun","run","index.ts"] run unprivileged; reference the Dockerfile's
WORKDIR, COPY steps, HEALTHCHECK, and CMD when applying the changes.
In `@storage/tests/e2e/local.rs`:
- Around line 22-37: The test currently calls spawn::spawn(...) and stores
handle but uses expect(...) on subsequent awaits (health::wait_for_healthy and
storage::backend::local::ensure_bucket) so a panic will skip cleaning up the
child; modify the test to ensure spawn::shutdown(handle) always runs on all
failure paths by converting the sequence to return a Result and using match/if
let Err(_) or a finally-like block (e.g., let res = (async { ...
health::wait_for_healthy(...).await?;
storage::backend::local::ensure_bucket(...).await?; Ok(()) }).await; if
res.is_err() { spawn::shutdown(handle).await.expect("shutdown"); } else {
spawn::shutdown(handle).await.expect("shutdown"); } ), or wrap handle in a guard
type that calls spawn::shutdown in Drop; apply the same pattern to the other
occurrence referenced (lines ~81-82) so the child process is always shut down.
In `@storage/tests/e2e/workers/harness/src/cases.ts`:
- Around line 7-15: The harness omits the new GCS backend: update the Provider
union type (Provider), the ALL_PROVIDERS array, and the BUCKETS map to include
'gcs' and a corresponding bucket name (e.g., 'scratch-gcs') so the TypeScript
harness will exercise GCS; also add a fake-gcs service to docker-compose.yml
(configure an appropriate emulator image/ports and ensure tests can reach it) so
the harness can run GCS-backed tests.
In `@storage/tests/integration.rs`:
- Around line 58-64: When spawning the worker Child fails, the earlier
background process named iii is not cleaned up; change the worker spawn handling
in boot() from using .spawn().ok()? to explicitly handle the Result and on Err
call iii.kill().ok() (or iii.kill().await/iii.wait().ok() as appropriate) before
returning None so the iii Child is terminated and no subprocess is leaked;
reference the worker Command::new(...).spawn() call and the iii Child variable
in your fix.
---
Minor comments:
In `@storage/README.md`:
- Around line 54-71: The example incorrectly uses iii.trigger with a
TriggerRequest to call the RPC handler storage::putObject; change the example to
use the Rust SDK's RPC/call/invoke method (the counterpart to TypeScript's
call('storage::putObject', ...)) instead of TriggerRequest/iii.trigger so the
RPC is invoked properly; locate usage of TriggerRequest and iii.trigger in the
README example and replace it with the SDK's call/invoke API, passing the
function id "storage::putObject" and the same payload and timeout parameters.
In `@storage/src/backend/gcs.rs`:
- Around line 200-219: The presign method currently always signs against
storage.googleapis.com and must early-fail when a custom endpoint is configured;
add a boolean field (e.g., endpoint_overridden) to the GcsBackend struct during
GcsBackend::build to record endpoint_url.is_some(), then modify async fn
presign(&self, req: PresignReq) to check that flag and return
Err(BackendError::PresignUnsupported("gcs presign always signs against
storage.googleapis.com; not supported when endpoint_url is overridden
(fake-gcs-server, etc.)".into())) when true, otherwise proceed with the existing
signing logic.
In `@storage/src/rustfs/health.rs`:
- Around line 20-27: The loop can exceed max_wait because each iteration calls
client.get(&url).send().await with a fixed request timeout; clamp each attempt
to the remaining time by computing let remaining =
deadline.saturating_duration_since(Instant::now()) and using either
client.get(&url).timeout(remaining.min(MAX_REQ_TIMEOUT)).send().await or
wrapping the send() in tokio::time::timeout(remaining.min(MAX_REQ_TIMEOUT),
client.get(&url).send()).await, then handle the timeout/error case the same as
other errors and continue/sleep; reference the existing symbols deadline,
Instant::now(), client.get(&url).send().await, and
sleep(Duration::from_millis(...)).await.
In `@storage/src/rustfs/spawn.rs`:
- Around line 174-181: The test function discover_binary_reads_env_var mutates
global process env with std::env::set_var/remove_var which races when tests run
in parallel; update the test (discover_binary_reads_env_var) to scope the
environment change using a thread-safe helper such as temp_env::with_var or
annotate/serialize the test using the serial_test crate (e.g., #[serial]) so
RUSTFS_BIN is set only for the test duration and automatically
restored/isolated, ensuring discover_binary() reads the temporary value without
global races.
In `@storage/src/triggers/pollers/pubsub.rs`:
- Around line 99-106: The current silent JSON-parse fallback turns malformed
pm.data into Null and causes normalize_gcs to emit a non-Ignored error that
leads to endless redelivery; replace the unwrap_or handling so a JSON parse Err
triggers a warning log (include pm.message_id and error info) and returns
NormalizeError::Ignored (so the dispatcher will ack), e.g. change the
serde_json::from_slice::<Value>(&pm.data).unwrap_or(...) usage in pubsub.rs to a
match/if let that logs warn and returns early with NormalizeError::Ignored
before constructing the envelope; this ensures normalize_gcs/dispatcher ack
poisoned messages instead of retrying forever.
In `@storage/src/triggers/pollers/sqs.rs`:
- Around line 94-104: The delete_message() call currently discards the Result
from send().await; change it to inspect the result and emit a warn! when it
fails (include the error, the receipt_handle value and self.queue_url for
context) while keeping the same non-panicking behavior; locate the block
handling msg.receipt_handle inside the poller (where delete_message(),
.queue_url(&self.queue_url) and .receipt_handle(rh) are invoked) and replace the
`let _ = ...send().await` with a match or if let Err(e) => warn!(...) that logs
e and identifying info.
In `@storage/tests/e2e/fixtures/minio-sqs-bridge/index.ts`:
- Around line 77-81: The code mapping body.Records into normalized assumes each
record is a non-null object and directly dereferences rec.eventName; update the
logic around normalized to first filter out non-object or null records (e.g.,
ensure typeof rec === 'object' && rec !== null) and only then cast to
Record<string, unknown> before reading eventName, and handle missing or
non-string eventName safely so primitives/nulls are skipped or returned
unchanged; change references in this block (body.Records and normalized)
accordingly.
In `@storage/tests/e2e/run-tests.sh`:
- Around line 335-353: The health-wait loop misclassifies services with no
healthcheck as unhealthy; update the awk filter in the unhealthy assignment (the
line that runs docker compose --profile cloud ps --format '{{.Service}}
{{.Health}}' and pipes to awk) so it ignores empty {{.Health}} fields (treat
empty as OK) — e.g. only print services whose second field is present and not
"healthy" — leaving the rest of the loop (health_deadline, break/exit logic, and
per-service logs) unchanged.
In `@storage/tests/e2e/s3.rs`:
- Around line 20-24: The test currently only early-returns when S3_TEST_BUCKET
is missing, causing unclear auth failures later; update the startup checks to
also require AWS_ACCESS_KEY_ID and AWS_SECRET_ACCESS_KEY (e.g., call require_env
for "AWS_ACCESS_KEY_ID" and "AWS_SECRET_ACCESS_KEY" along with "S3_TEST_BUCKET"
before constructing BucketConfig/S3BucketConfig) and return early if either
credential is unset so the test is skipped with a clear reason.
In `@storage/tests/e2e/workers/harness/src/cases-concurrency.ts`:
- Around line 17-19: The current assertion using etags.size >= 1 is too weak and
can pass when many put responses are missing etag; update the check around the
etags Set created from puts.map(...) and the assertTruthy call to verify every
put has a defined, non-empty etag (e.g., assert that puts.every(p => p.etag !=
null && p.etag !== '') and that etags.size === puts.length) so missing/undefined
etag fields fail the test; refer to the etags variable, the puts array mapping,
and the assertTruthy invocation to locate and update the assertion.
In `@storage/tests/e2e/workers/harness/src/cases-provider.ts`:
- Around line 17-35: The test's presign URL assertion only checks
r.url.includes(ctx.bucket) but must actually validate URL parseability and that
the bucket appears either in hostname or pathname; update the run(ctx) for the
'[s3] presigned URL is path-style against MinIO' case to attempt parsing r.url
with the URL constructor (ensuring it doesn't throw), then assert that either
parsed.hostname or parsed.pathname contains ctx.bucket (so both virtual-host and
path-style URLs pass), using the existing r from the storage::presignUrl call to
produce clear assertion messages if parsing fails or the bucket is not found.
In `@storage/tests/e2e/workers/harness/src/runner.ts`:
- Around line 181-187: resetEvents currently clears waiter timers and drops
waiter entries causing any in-flight waitForEvent promises to hang; modify
resetEvents (which manipulates this.waiters, this.events and each waiter.timer)
to first reject each pending waiter (call its rejector with a clear Error/Cancel
reason), then clear its timer, and finally clear this.waiters and this.events so
no promise is left unresolved; ensure you reference the waiter objects used by
waitForEvent so their resolve/reject callbacks are invoked before removal.
In `@storage/tests/e2e/workers/harness/src/worker.ts`:
- Around line 17-27: The code blindly casts environment strings to Provider[]
via PROVIDERS_RAW/PROVIDERS and TRIGGER_PROVIDERS_RAW/TRIGGER_PROVIDERS which
can hide invalid provider names; update parsing to validate each token against
the canonical set of allowed Provider values (e.g., an enum or array of valid
provider strings) rather than using "as Provider[]", and if any token is not
recognized throw a clear, early error (with the invalid names listed) during
initialization; factor this into a small helper like parseProviders(raw:
string): Provider[] and use it for both PROVIDERS and TRIGGER_PROVIDERS to
ensure fail-fast validation.
---
Nitpick comments:
In @.github/workflows/storage-e2e.yml:
- Around line 64-68: The workflow step "Install iii engine" currently pulls
install.sh from the main branch which causes non-reproducible CI; update that
step to pin the installer to a specific, committed release/tag or version (e.g.,
replace the "main" URL with a release/tag or add an explicit version
parameter/ENV used by the installer) so every run installs the same iii engine
version; ensure the step name "Install iii engine" and the curl invocation are
updated accordingly and include a comment or variable (e.g., III_VERSION) so
future updates require an explicit version bump.
In `@storage/build.rs`:
- Line 4: Replace the unwrap() call on std::env::var("TARGET") in build.rs with
expect() to provide a clearer, actionable error message; update the expression
where std::env::var("TARGET").unwrap() appears so it uses expect("TARGET must be
set by Cargo build scripts") (or similar) to document intent and improve the
build-time panic message.
In `@storage/src/triggers/dispatcher.rs`:
- Around line 36-111: Sequential per-subscriber awaits in dispatch cause total
latency to sum per-subscriber time; change dispatch in
storage::triggers::dispatcher.rs to fire all per-subscriber trigger futures
concurrently (e.g., map subs from registry.subscribers_for(...) to individual
futures that call self.iii.trigger(req) wrapped with the same timeout logic
using sub.handler_timeout_ms) and then await them together with
futures::future::join_all or similar, preserving the existing ack extraction
logic (null => true, value.get("ack")... ) and keeping function_id and
timeout_ms logging behavior.
In `@storage/src/triggers/handler.rs`:
- Around line 15-98: The two handlers ObjectCreatedHandler and
ObjectDeletedHandler duplicate register_trigger/unregister_trigger logic;
consolidate into one generic handler to avoid drift. Replace both with a single
struct (e.g., ObjectEventHandler) parameterized by EventKind and a parser
function/enum specifying which config type to deserialize (CreatedConfig vs
DeletedConfig) so register_trigger calls the chosen parser, validates
wired_buckets, calls registry.register with the injected EventKind, and emits a
log message derived from that EventKind; keep unregister_trigger delegating to
registry.unregister(&config.id) unchanged. Implement the parser injection as a
Fn(&Value) -> Result<ParsedConfig, IIIError> or an enum branch inside the
generic handler to locate logic around register_trigger in
ObjectCreatedHandler/ObjectDeletedHandler and switch to the unified
register_trigger implementation.
In `@storage/src/triggers/pollers/rustfs_webhook.rs`:
- Around line 32-68: spawn_receiver duplicates spawn_receiver_on except for the
bind argument; change spawn_receiver to simply call and return
spawn_receiver_on(state, 0) so the ephemeral-port behavior is delegated to
spawn_receiver_on, eliminating the duplicated router/serve/channels code (keep
the same async signature and return type for spawn_receiver and remove the
duplicate implementation).
In `@storage/tests/e2e/run-tests.sh`:
- Around line 253-272: The current SHA256 verification silently skips when the
SHA256SUMS file is reachable but awk finds no entry for "$asset"; change this to
fail when the caller requires strict verification by checking an env var
(RUSTFS_REQUIRE_SHA256=1): in the run-tests.sh block that computes
expected="$(awk -v a="$asset" '$2==a {print $1}' "$tmpdir/SHA256SUMS")", if [[
-z "$expected" ]] then if RUSTFS_REQUIRE_SHA256 is set to 1 exit with an error
(return 1) and log a clear fatal message referencing "$asset" and
"$tmpdir/SHA256SUMS", otherwise keep the existing warn-and-skip behavior;
preserve the existing sha256sum/shasum branch and subsequent checks for actual
vs expected.
In `@storage/tests/e2e/workers/harness/src/cases.ts`:
- Around line 163-168: The assertion message for the delete check is misleading:
update the assertion in the block that calls ctx.call('storage::deleteObject')
so the assertEqual on del.deleted uses a positive expectation message (e.g.,
"expected del.deleted=true" or "expected delete.deleted to be true") instead of
"delete returned deleted=false"; target the assertEqual call that compares
del.deleted to true to make the message reflect the expected state.
- Line 41: The JSDoc above the call property contains a typo "an storage" —
update the comment to read "a storage" (or reword to "the storage" if more
appropriate) for the call: (functionId: string, payload: unknown) =>
Promise<any>; identifier so the documentation is grammatically correct.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
…ache Two failures in the first storage-e2e run on PR #91: 1. Script self-tests job: --filter and SIGINT tests bailed early with "iii engine binary missing" because the job didn't install iii. Add the same `Install iii engine` step the harness job already has — script-tests/run.sh exercises run-tests.sh's full preflight path, so the engine has to be present even though we never actually drive it. 2. Harness (--providers=all) job: actions/setup-node@v4 hard-failed with "Some specified paths were not resolved, unable to cache dependencies" because cache: 'npm' was paired with a cache-dependency-path pointing at a package-lock.json we don't commit (run-tests.sh runs `npm install` at first run, generating one locally). Drop the cache hint; ~10s slower per run, but unblocks the job. Revisit if we ever commit a lockfile (mirror the iii-database/shell pattern).
Correctness:
- backend/factory: pipe S3 force_path_style from config (was hardcoded false)
- backend/local: only call create_bucket on definite 404; surface other
head_bucket errors instead of masking them
- backend/local: propagate create_dir_all errors for queue_dir
- backend/s3: fail closed on GET when Content-Length missing and
max_inline_bytes set (closes cap bypass)
- backend/gcs: presign returns PresignUnsupported when endpoint_url is
overridden (signing always targets storage.googleapis.com)
- handlers/put_object: accept legacy "body" via serde alias; reject
oversized base64 payloads before allocating decode buffer
- triggers/handler: build CONFIG_ERROR envelopes via serde_json::json!
so embedded errors and bucket names are properly escaped
- triggers/normalize: decode '+' to space before percent-decoding S3 keys
- triggers/pollers/cf_queue: ack and log poison messages on normalize
errors instead of redelivering forever
- triggers/pollers/pubsub: ack messages with non-JSON data instead of
silently coercing to Null and looping forever
- triggers/pollers/sqs: log delete_message failures with queue context
Tests / fixtures:
- tests/integration: kill iii Child if worker spawn fails (prevent leak)
- tests/e2e/local: wrap test body so spawn::shutdown always runs
- tests/e2e/s3: skip cleanly when AWS credentials are missing
- tests/e2e/run-tests.sh: treat empty {{.Health}} field as ok
- fixtures/minio-sqs-bridge/Dockerfile: run as non-root bun user
- fixtures/minio-sqs-bridge/index.ts: filter null/non-object records
- harness/cases-concurrency: assert every put returned a non-empty etag
- harness/cases-provider: parse presigned URL and check bucket in
hostname or pathname
- harness/runner: reject pending waiters in resetEvents so awaits
don't hang
- harness/worker: validate HARNESS_PROVIDERS / HARNESS_TRIGGER_PROVIDERS
against ALL_PROVIDERS instead of blind cast
Misc:
- build.rs: clearer panic message via expect()
- harness/cases: fix typo and improve delete assertion message
Summary
```
80 files, +13,701 LOC. Net new code; no existing files touched outside CI + .gitignore.
```
Pre-landing review (from `/review`)
No CRITICAL findings. Six items surfaced for follow-up — none block landing on their own, but at least the first two should be addressed before any production-ish deploy:
Full review notes are in `~/.gstack/projects/iii-hq-workers/featiii-storage-worker-reviews.jsonl`.
Notes
Test plan
Summary by CodeRabbit
New Features
Documentation
Testing