summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--server/src/worker.rs149
-rw-r--r--tests/worker/test_data.py52
-rw-r--r--worker/data.py30
3 files changed, 205 insertions, 26 deletions
diff --git a/server/src/worker.rs b/server/src/worker.rs
index 8e28de5..d046277 100644
--- a/server/src/worker.rs
+++ b/server/src/worker.rs
@@ -194,31 +194,73 @@ pub async fn docker_available() -> bool {
.unwrap_or(false)
}
-/// Mounts need world permissions: the container runs as uid 65534 while host
-/// ownership is the server user. Best effort only.
-fn prepare_mounts(mounts: &[(String, String, bool)]) {
- for (src, _dst, ro) in mounts {
+/// Mounts need world permissions for the non-root container (uid 65534).
+///
+/// Strict scope (leader review 2026-09-17):
+/// - standalone FILE mounts are touched ONLY when they are read-only mount
+/// of a worker data feed (dst starts with `/data/`) — never arbitrary
+/// host files; symlinks are rejected, chmod errors surface, nothing is
+/// silently swallowed.
+/// - directory mounts (job input/output dirs) still normalize perms; any
+/// failure is now returned instead of swallowed.
+fn prepare_mounts(mounts: &[(String, String, bool)]) -> AppResult<()> {
+ for (src, dst, ro) in mounts {
let p = std::path::Path::new(src);
+ // Standalone file mounts (e.g. dataset objects bound directly to
+ // /data/<name> for backtests) must independently become world
+ // readable; they do not live under a prepare_mounts ro directory.
if !p.is_dir() {
+ let md = std::fs::symlink_metadata(p)
+ .map_err(|e| crate::error::AppError::internal(format!("mount {src:?} unreachable: {e}")))?;
+ if md.file_type().is_symlink() {
+ return Err(crate::error::AppError::internal(format!(
+ "symlinked mount rejected: {src}"
+ )));
+ }
+ let data_ro = *ro && dst.starts_with("/data/");
+ if data_ro {
+ if !p.is_file() {
+ return Err(crate::error::AppError::internal(format!(
+ "read-only data mount is not a regular file: {src}"
+ )));
+ }
+ std::fs::set_permissions(p, std::fs::Permissions::from_mode(0o644)).map_err(
+ |e| {
+ crate::error::AppError::internal(format!(
+ "cannot make data mount readable: {src}: {e}"
+ ))
+ },
+ )?;
+ }
+ // Anything else (non-/data or non-ro) is intentionally untouched.
continue;
}
let mode = if *ro { 0o755 } else { 0o777 };
- let _ = std::fs::set_permissions(p, std::fs::Permissions::from_mode(mode));
+ std::fs::set_permissions(p, std::fs::Permissions::from_mode(mode)).map_err(|e| {
+ crate::error::AppError::internal(format!("cannot set dir mount permissions: {src}: {e}"))
+ })?;
// Files inside ro input dirs must be world readable; output files are
// written by the container with its umask.
if *ro {
- if let Ok(rd) = std::fs::read_dir(p) {
- for e in rd.flatten() {
- let fmode = if e.path().is_file() {
- std::fs::Permissions::from_mode(0o644)
- } else {
- std::fs::Permissions::from_mode(0o755)
- };
- let _ = std::fs::set_permissions(e.path(), fmode);
- }
+ for e in std::fs::read_dir(p)
+ .map_err(|e| crate::error::AppError::internal(format!("cannot read mount dir {src}: {e}")))?
+ .flatten()
+ {
+ let fmode = if e.path().is_file() {
+ std::fs::Permissions::from_mode(0o644)
+ } else {
+ std::fs::Permissions::from_mode(0o755)
+ };
+ std::fs::set_permissions(e.path(), fmode).map_err(|err| {
+ crate::error::AppError::internal(format!(
+ "cannot fix ro mount entry {:?}: {err}",
+ e.path()
+ ))
+ })?;
}
}
}
+ Ok(())
}
/// Run the worker image with a fixed container name so cancel maps to one
@@ -232,7 +274,9 @@ pub async fn run_named(
timeout_secs: u64,
) -> AppResult<ContainerResult> {
let cfg = &cx.cfg;
- prepare_mounts(mounts);
+ prepare_mounts(mounts).map_err(|e| {
+ crate::error::AppError::internal(format!("worker mount preparation failed: {}", e.message))
+ })?;
let full = docker_args(network, mounts, &cfg.worker_image, name, args);
execute_docker(&full, name, timeout_secs).await
}
@@ -480,3 +524,78 @@ mod tests {
assert_eq!(truncate("short", 100), "short");
}
}
+
+
+
+#[cfg(test)]
+mod mount_perm_tests {
+ use super::prepare_mounts;
+ use std::os::unix::fs::PermissionsExt;
+
+ fn mode_of(p: &std::path::Path) -> u32 {
+ std::fs::metadata(p).unwrap().permissions().mode() & 0o777
+ }
+
+ /// 159399 candidate QA regression (2026-09-17): a run failed because the
+ /// directly-mounted dataset object file kept the server process' umask
+ /// perms (0o600) while the container runs as uid 65534 =>
+ /// PermissionError in the backtest worker. prepare_mounts already fixed
+ /// files *inside* ro directories but skipped standalone file mounts.
+ #[test]
+ fn standalone_data_file_mounts_become_world_readable() {
+ let td = tempfile::tempdir().unwrap();
+ let f = td.path().join("48c.csv");
+ std::fs::write(&f, b"date,open\n2026-01-02,1.0\n").unwrap();
+ std::fs::set_permissions(&f, std::fs::Permissions::from_mode(0o600)).unwrap();
+ prepare_mounts(&[(f.display().to_string(), "/data/48c.csv".into(), true)]).unwrap();
+ assert_eq!(mode_of(&f) & 0o004, 0o004, "others-read must be set on data mounts");
+ }
+
+ /// Narrow scope: file mounts that are NOT read-only /data feeds must be
+ /// left exactly as they are (no blanket chmod of arbitrary paths).
+ #[test]
+ fn non_data_file_mounts_are_untouched() {
+ let td = tempfile::tempdir().unwrap();
+ for (dst, ro) in [("/input/f.json", false), ("/data/x", false)] {
+ let f = td.path().join(format!("f{}.bin", dst.replace('/', "_")));
+ std::fs::write(&f, b"x").unwrap();
+ std::fs::set_permissions(&f, std::fs::Permissions::from_mode(0o600)).unwrap();
+ prepare_mounts(&[(f.display().to_string(), dst.into(), ro)]).unwrap();
+ assert_eq!(mode_of(&f) & 0o077, 0, "non-data file mount must be untouched: {}", dst);
+ }
+ }
+
+ /// Symlinks are rejected with an explicit error; the link target must not
+ /// be widened (no chmod via a link).
+ #[test]
+ fn symlinked_mount_is_rejected_not_chmodded() {
+ let td = tempfile::tempdir().unwrap();
+ let target = td.path().join("real-data.csv");
+ std::fs::write(&target, b"date\n2026-01-02\n").unwrap();
+ std::fs::set_permissions(&target, std::fs::Permissions::from_mode(0o600)).unwrap();
+ let link = td.path().join("data.csv");
+ std::os::unix::fs::symlink(&target, &link).unwrap();
+ let err = prepare_mounts(&[(
+ link.display().to_string(),
+ "/data/data.csv".into(),
+ true,
+ )])
+ .expect_err("symlink must be rejected");
+ assert_eq!(mode_of(&target) & 0o077, 0o000, "target must stay untouched");
+ assert!(err.message.contains("symlink"), "{:?}", err.message);
+ }
+
+ #[test]
+ fn directory_mounts_keep_fixing_files_inside_ro() {
+ let td = tempfile::tempdir().unwrap();
+ let d = td.path().join("input");
+ std::fs::create_dir_all(&d).unwrap();
+ let f = d.join("request.json");
+ std::fs::write(&f, b"{}").unwrap();
+ std::fs::set_permissions(&d, std::fs::Permissions::from_mode(0o700)).unwrap();
+ std::fs::set_permissions(&f, std::fs::Permissions::from_mode(0o600)).unwrap();
+ prepare_mounts(&[(d.display().to_string(), "/input".into(), true)]).unwrap();
+ assert_eq!(mode_of(&d) & 0o044, 0o044, "ro dir needs traversal");
+ assert_eq!(mode_of(&f) & 0o004, 0o004, "files inside ro dir need others-read");
+ }
+}
diff --git a/tests/worker/test_data.py b/tests/worker/test_data.py
index d6eeae1..e5c30a3 100644
--- a/tests/worker/test_data.py
+++ b/tests/worker/test_data.py
@@ -334,3 +334,55 @@ def test_fetch_columns_subset_requested():
assert endpoint == "stock_zh_a_hist"
assert all(c in df.columns for c in ("date", "symbol", "open", "close", "volume"))
monkeypatch.undo()
+
+
+# ---- identity contract: the server persists instruments with the user-facing
+# market field (SH/SZ/BJ) exactly as the frontend manual-entry form sends them
+# (159399 QA regression 2026-09-17: market:"SZ" used to fail with
+# unsupported_market before any network attempt). Canonical "SH#600000" and
+# market:"cn" bare-code inference remain the primary, unchanged contracts. ----
+
+def test_split_identity_accepts_bare_code_with_market_sz_matches_device_inference():
+ exch, code = data.split_identity({"symbol": "159399", "market": "SZ"})
+ assert (exch, code) == ("SZ", "159399")
+
+
+def test_split_identity_accepts_bare_code_with_market_sh():
+ exch, code = data.split_identity({"symbol": "600000", "market": "SH"})
+ assert (exch, code) == ("SH", "600000")
+
+
+def test_split_identity_accepts_bare_code_with_market_cn_inference_unchanged():
+ exch, code = data.split_identity({"symbol": "159399", "market": "cn"})
+ assert (exch, code) == ("SZ", "159399")
+
+
+def test_etf_fetch_passes_market_sz_instrument_to_provider(monkeypatch):
+ """Full fetch_source identity path: market:"SZ" (as stored by the server
+ from the frontend manual-entry form) must reach the sina provider."""
+ captured = {}
+ real_frame = _sina_synth_frame()
+
+ def fake_em(inst, s, e, f, adj):
+ raise data.DataError("source_unavailable", "eastmoney down (test)")
+
+ def fake_sina(inst, s, e, f, adj):
+ captured["instrument"] = inst
+ return data.FetchResult(real_frame.copy(), "fund_etf_hist_sina",
+ {"symbol": "sz159399"}, real_frame.copy(), "sina")
+
+ monkeypatch.setattr(data, "_fetch_eastmoney", fake_em)
+ monkeypatch.setattr(data, "_fetch_sina", fake_sina)
+ inst = {"symbol": "159399", "market": "SZ", "asset_type": "etf", "name": "国泰自由现金流"}
+ res, warns = data.fetch_source_with_warnings(inst, "2025-12-31", "2026-09-17",
+ "daily", "none",
+ ["open", "high", "low", "close", "volume"],
+ source="auto")
+ assert captured["instrument"] is inst
+ assert res.provider == "sina"
+ assert any(w.startswith("provider_fallback") for w in warns)
+
+
+def test_split_identity_rejects_unknown_market_prefixes():
+ with pytest.raises(data.DataError):
+ data.split_identity({"symbol": "1", "market": "adlhkj"})
diff --git a/worker/data.py b/worker/data.py
index 841a1eb..13f8808 100644
--- a/worker/data.py
+++ b/worker/data.py
@@ -66,13 +66,20 @@ def split_identity(instrument: dict) -> tuple[str, str]:
raise DataError("bad_identity", f"unknown exchange prefix: {exch}")
else:
code = symbol
- market = instrument.get("market", "cn")
- if market != "cn":
- raise DataError("unsupported_market", f"market not supported by this adapter: {market}")
- if len(code) == 6 and code[0] in "369":
- exch = "SH"
+ market = str(instrument.get("market", "cn")).strip() or "cn"
+ # The server persists the user-facing market field verbatim (SH/SZ/BJ
+ # from the manual-entry form); "cn" (search-sourced items, any case)
+ # keeps the original code-prefix inference. Unknown markets must still
+ # be rejected — never guessed.
+ if market.upper() in IDENTITY_EXCHANGES:
+ exch = market.upper()
+ elif market.lower() == "cn":
+ if len(code) == 6 and code[0] in "369":
+ exch = "SH"
+ else:
+ exch = "SZ"
else:
- exch = "SZ"
+ raise DataError("unsupported_market", f"market not supported by this adapter: {market}")
if not code.isdigit() or len(code) != 6:
raise DataError("bad_identity", "A-share code must be a 6-digit number")
return exch, code
@@ -157,11 +164,12 @@ def _fetch_sina(instrument: dict, start: str, end: str,
# requested window (ISO strings compare lexicographically like dates)
df = df[(df["date"] >= start) & (df["date"] <= end)].reset_index(drop=True)
df.attrs["source_warnings"] = [
- "sina returns the full trading history without date parameters; sliced "
- "locally to the requested range; data is unadjusted (no adjust parameter)",
- "sina volume unit is 股 (shares); verified to be 100x the 手 (lots) "
- "convention used by lot-based providers on the same session "
- "(2026-09-17 evidence, docs/recovery-01-plan.md)",
+ "sina fund_etf_hist_sina has no date-range parameter: the provider "
+ "returned its available history and it was sliced locally to the "
+ "requested range (actual coverage differences are reported separately)",
+ "sina ETF klines are unadjusted: the interface has no adjust parameter",
+ "sina fund_etf_hist_sina 成交量单位为股;数值按 provider 原样保留,未做换算或取整。"
+ "与其他行情接口对比时,请先核对各 provider 自己公布的口径,不要按惯例直接换算",
]
return FetchResult(df, endpoint, params, raw, "sina")