修复 异常BUG
This commit is contained in:
@@ -800,7 +800,12 @@ impl 线段 {
|
||||
|
||||
seg.序号
|
||||
.store(之前线段.序号.load(Ordering::Relaxed) + 1, Ordering::Relaxed);
|
||||
*seg.前一缺口.write().unwrap() = Self::获取缺口(之前线段);
|
||||
*seg.前一缺口.write().unwrap() = if 之前线段.短路修正.load(Ordering::Relaxed)
|
||||
{
|
||||
None
|
||||
} else {
|
||||
Self::获取缺口(之前线段)
|
||||
};
|
||||
*seg.前一结束位置.write().unwrap() = Some(Arc::clone(
|
||||
之前线段.基础序列.read().unwrap().last().unwrap(),
|
||||
));
|
||||
@@ -970,6 +975,7 @@ impl 线段 {
|
||||
}
|
||||
|
||||
let 左 = Arc::clone(&基础序列[基础序列.len() - 3]);
|
||||
let 中 = Arc::clone(&基础序列[基础序列.len() - 2]);
|
||||
let 右 = Arc::clone(&基础序列[基础序列.len() - 1]);
|
||||
|
||||
// 方向条件
|
||||
@@ -987,7 +993,8 @@ impl 线段 {
|
||||
基础序列
|
||||
);
|
||||
|
||||
let 原始基础序列 = 当前线段.基础序列.read().unwrap().clone();
|
||||
// Reassign to full copy (matching Python pattern)
|
||||
let 基础序列 = 当前线段.基础序列.read().unwrap().clone();
|
||||
Self::_弹出线段(
|
||||
线段序列,
|
||||
&Arc::clone(线段序列.last().unwrap()),
|
||||
@@ -1021,13 +1028,13 @@ impl 线段 {
|
||||
cur.特征序列.write().unwrap()[2] = None;
|
||||
|
||||
let 开始笔 = Arc::clone(cur.基础序列.read().unwrap().last().unwrap());
|
||||
let 开始序号 = 原始基础序列
|
||||
let 开始序号 = 基础序列
|
||||
.iter()
|
||||
.position(|x| Arc::as_ptr(x) == Arc::as_ptr(&开始笔));
|
||||
|
||||
开始序号_opt = 开始序号;
|
||||
if let Some(序号) = 开始序号 {
|
||||
待添加元素 = 原始基础序列[序号 + 1..].to_vec();
|
||||
待添加元素 = 基础序列[序号 + 1..].to_vec();
|
||||
} else {
|
||||
待添加元素 = Vec::new();
|
||||
}
|
||||
@@ -1047,19 +1054,14 @@ impl 线段 {
|
||||
let 当前线段 = Arc::clone(&线段序列[idx]);
|
||||
当前线段.短路修正.store(true, Ordering::Relaxed);
|
||||
if 当前线段.特征序列.read().unwrap()[2].is_some() {
|
||||
let 段 = 虚线::创建线段(&[
|
||||
Arc::clone(&基础序列[基础序列.len() - 3]),
|
||||
Arc::clone(&基础序列[基础序列.len() - 2]),
|
||||
Arc::clone(&基础序列[基础序列.len() - 1]),
|
||||
]);
|
||||
let 段 = 虚线::创建线段(&[Arc::clone(&左), Arc::clone(&中), Arc::clone(&右)]);
|
||||
let 段_rc = Arc::new(段);
|
||||
Self::_添加线段(线段序列, 段_rc, 配置, format!("{}, {}", line!(), 层级));
|
||||
|
||||
// Set feature sequence [0]
|
||||
let 新段 = Self::取段(线段序列.last_mut().unwrap());
|
||||
let 中笔 = Arc::clone(&基础序列[基础序列.len() - 2]);
|
||||
新段.特征序列.write().unwrap()[0] =
|
||||
Some(Arc::new(线段特征::新建(vec![中笔], 新段.方向())));
|
||||
Some(Arc::new(线段特征::新建(vec![中], 新段.方向())));
|
||||
}
|
||||
|
||||
true
|
||||
|
||||
@@ -32,18 +32,15 @@ use std::sync::RwLock;
|
||||
use tracing::{error, info};
|
||||
|
||||
/// 立体分析器 — 多周期协调器
|
||||
///
|
||||
/// 包含一个K线合成器和每周期一个观察者。
|
||||
/// 输入最小周期K线,合成大周期后分发到对应观察者。
|
||||
pub struct 立体分析器 {
|
||||
pub 周期组: Vec<i64>,
|
||||
输入周期: i64,
|
||||
K线合成器: K线合成器,
|
||||
单体分析器: HashMap<i64, Arc<RwLock<观察者>>>,
|
||||
pub K线合成器: K线合成器,
|
||||
pub 单体分析器: HashMap<i64, Arc<RwLock<观察者>>>,
|
||||
}
|
||||
|
||||
impl 立体分析器 {
|
||||
/// 创建立体分析器,自动创建K线合成器 + 每周期一个观察者
|
||||
/// 创建立体分析器 — 对应 Python 立体分析器.__init__
|
||||
pub fn new(
|
||||
符号: String,
|
||||
周期组: Vec<i64>,
|
||||
@@ -58,9 +55,7 @@ impl 立体分析器 {
|
||||
let 默认配置 = 配置.unwrap_or_default();
|
||||
let 配置组 = 配置组.unwrap_or_default();
|
||||
|
||||
let K线合成器 = K线合成器::new(符号.clone(), 周期组.clone());
|
||||
|
||||
let mut 单体分析器 = HashMap::new();
|
||||
let mut 单体分析器: HashMap<i64, Arc<RwLock<观察者>>> = HashMap::new();
|
||||
for &周期 in &周期组 {
|
||||
let mut 当前配置 = 配置组
|
||||
.get(&周期)
|
||||
@@ -101,6 +96,18 @@ impl 立体分析器 {
|
||||
}
|
||||
}
|
||||
|
||||
// 对应 Python: K线合成器(符号, 周期组, self.__K线回调)
|
||||
let 单体分析器_回调 = 单体分析器.clone();
|
||||
let K线合成器 = K线合成器::new(
|
||||
符号.clone(),
|
||||
周期组.clone(),
|
||||
Some(Box::new(
|
||||
move |_信号类型: String, _标识: String, 周期: i64, 完成K线: K线| {
|
||||
立体分析器::__K线回调_调度(&单体分析器_回调, 周期, 完成K线);
|
||||
},
|
||||
)),
|
||||
);
|
||||
|
||||
Self {
|
||||
周期组,
|
||||
输入周期,
|
||||
@@ -109,8 +116,28 @@ impl 立体分析器 {
|
||||
}
|
||||
}
|
||||
|
||||
/// 投喂K线 — 统一入口,接收最小周期K线
|
||||
/// 匹配 Python __K线回调:合成器完成K线时喂给观察者
|
||||
/// __K线回调 — 对应 Python 立体分析器.__K线回调
|
||||
fn __K线回调(&self, _信号类型: String, _标识: String, 周期: i64, 完成K线: K线) {
|
||||
if let Some(观察员) = self.单体分析器.get(&周期) {
|
||||
let mut obs = 观察员.write().unwrap();
|
||||
obs.增加原始K线(完成K线);
|
||||
// 对应 Python: if 当前K线 := self._K线合成器.获取当前K线(周期)
|
||||
// _完成K线刚清空当前K线,获取当前K线返回 None,所以这里不添加
|
||||
}
|
||||
}
|
||||
|
||||
/// 静态调度版本 — 用于回调闭包
|
||||
fn __K线回调_调度(
|
||||
单体分析器: &HashMap<i64, Arc<RwLock<观察者>>>,
|
||||
周期: i64,
|
||||
完成K线: K线,
|
||||
) {
|
||||
if let Some(观察员) = 单体分析器.get(&周期) {
|
||||
观察员.write().unwrap().增加原始K线(完成K线);
|
||||
}
|
||||
}
|
||||
|
||||
/// 投喂K线 — 对应 Python 立体分析器.投喂K线
|
||||
pub fn 投喂K线(&mut self, 普K: K线) {
|
||||
if 普K.周期 != self.输入周期 {
|
||||
panic!(
|
||||
@@ -118,19 +145,7 @@ impl 立体分析器 {
|
||||
普K.周期, self.输入周期
|
||||
);
|
||||
}
|
||||
|
||||
// Feed to synthesizer, get completion events
|
||||
let 完成事件 = self.K线合成器.投喂K线(普K);
|
||||
|
||||
// Dispatch on completion events (matching Python's __K线回调)
|
||||
for (周期, 完成K线) in 完成事件 {
|
||||
if let Some(观察员) = self.单体分析器.get(&周期) {
|
||||
观察员.write().unwrap().增加原始K线(完成K线);
|
||||
if let Some(当前K线) = self.K线合成器.获取当前K线(周期) {
|
||||
观察员.write().unwrap().增加原始K线(当前K线.clone());
|
||||
}
|
||||
}
|
||||
}
|
||||
self.K线合成器.投喂K线(普K);
|
||||
}
|
||||
|
||||
/// 获取指定周期的观察者
|
||||
@@ -138,8 +153,7 @@ impl 立体分析器 {
|
||||
self.单体分析器.get(&周期).cloned()
|
||||
}
|
||||
|
||||
/// 测试_保存数据 — 多级别数据拆分保存
|
||||
/// 创建父目录 PyM_{标识}_{起始时间}_{结束时间},各周期观察者保存到子目录
|
||||
/// 测试_保存数据 — 对应 Python 立体分析器.测试_保存数据
|
||||
pub fn 测试_保存数据(&self, root: Option<&str>) {
|
||||
let 根目录 = match root {
|
||||
Some(r) => std::path::PathBuf::from(r),
|
||||
@@ -163,7 +177,6 @@ impl 立体分析器 {
|
||||
.get(&self.输入周期)
|
||||
.map(|o| o.read().unwrap().符号.clone())
|
||||
.unwrap_or_default();
|
||||
|
||||
let 周期 = self
|
||||
.单体分析器
|
||||
.get(&self.输入周期)
|
||||
@@ -189,4 +202,30 @@ impl 立体分析器 {
|
||||
|
||||
info!("多级别数据拆分保存完成,目录:{}", 保存路径.display());
|
||||
}
|
||||
|
||||
/// 相等 — 各周期观察者全量比对,对应 Python `立体分析器相等`
|
||||
pub fn 相等(&self, other: &Self, 浮点容差: f64) -> (bool, String) {
|
||||
let 标签 = format!("立体分析器校验[A={:?},B={:?}]", self.周期组, other.周期组);
|
||||
|
||||
if self.周期组 != other.周期组 {
|
||||
return (false, format!("{标签}: 周期组不一致"));
|
||||
}
|
||||
|
||||
for 周期 in &self.周期组 {
|
||||
let a_obs = match self.单体分析器.get(周期) {
|
||||
Some(o) => o.read().unwrap(),
|
||||
None => return (false, format!("{标签}: 周期{周期} 观察者不存在 (A)")),
|
||||
};
|
||||
let b_obs = match other.单体分析器.get(周期) {
|
||||
Some(o) => o.read().unwrap(),
|
||||
None => return (false, format!("{标签}: 周期{周期} 观察者不存在 (B)")),
|
||||
};
|
||||
let (eq, msg) = a_obs.相等(&b_obs, 浮点容差);
|
||||
if !eq {
|
||||
return (false, format!("{标签}: 周期{周期} >> {msg}"));
|
||||
}
|
||||
}
|
||||
|
||||
(true, format!("{标签}:所有周期观察者全量校验全部一致"))
|
||||
}
|
||||
}
|
||||
|
||||
+260
-114
@@ -261,78 +261,103 @@ impl 观察者 {
|
||||
None => return,
|
||||
};
|
||||
|
||||
// Step 2: 笔分析(无条件)
|
||||
笔::分析(
|
||||
当前分型,
|
||||
&mut self.分型序列,
|
||||
&mut self.笔序列,
|
||||
&self.缠论K线序列,
|
||||
&self.普通K线序列,
|
||||
0,
|
||||
&self.配置,
|
||||
);
|
||||
// Step 2: 笔分析
|
||||
if self.配置.分析笔 {
|
||||
笔::分析(
|
||||
当前分型,
|
||||
&mut self.分型序列,
|
||||
&mut self.笔序列,
|
||||
&self.缠论K线序列,
|
||||
&self.普通K线序列,
|
||||
0,
|
||||
&self.配置,
|
||||
);
|
||||
}
|
||||
if self.分型序列.is_empty() {
|
||||
return;
|
||||
}
|
||||
|
||||
// Step 3: 笔中枢分析(无条件)
|
||||
中枢::分析(&self.笔序列, &mut self.笔_中枢序列, true, "", 0);
|
||||
// Step 3: 笔中枢分析
|
||||
if self.配置.分析笔中枢 {
|
||||
中枢::分析(&self.笔序列, &mut self.笔_中枢序列, true, "", 0);
|
||||
}
|
||||
if self.笔序列.is_empty() {
|
||||
return;
|
||||
}
|
||||
|
||||
// Step 4: 线段分析 — 3 级递归
|
||||
for i in 0..self.线段分析层次 {
|
||||
if i == 0 {
|
||||
线段::分析(
|
||||
&self.笔序列,
|
||||
&mut self.线段序列组[i],
|
||||
&self.配置,
|
||||
0,
|
||||
&[相对方向::向上, 相对方向::向下],
|
||||
);
|
||||
} else {
|
||||
let 源序列 = self.线段序列组[i - 1].clone();
|
||||
线段::分析(
|
||||
&源序列,
|
||||
&mut self.线段序列组[i],
|
||||
&self.配置,
|
||||
0,
|
||||
&[相对方向::向上, 相对方向::向下],
|
||||
);
|
||||
if self.配置.分析线段 || self.配置.分析线段中枢 {
|
||||
for i in 0..self.线段分析层次 {
|
||||
if i == 0 {
|
||||
if self.配置.分析线段 {
|
||||
线段::分析(
|
||||
&self.笔序列,
|
||||
&mut self.线段序列组[i],
|
||||
&self.配置,
|
||||
0,
|
||||
&[相对方向::向上, 相对方向::向下],
|
||||
);
|
||||
}
|
||||
} else {
|
||||
if self.配置.分析线段 {
|
||||
let 源序列 = self.线段序列组[i - 1].clone();
|
||||
线段::分析(
|
||||
&源序列,
|
||||
&mut self.线段序列组[i],
|
||||
&self.配置,
|
||||
0,
|
||||
&[相对方向::向上, 相对方向::向下],
|
||||
);
|
||||
}
|
||||
}
|
||||
if self.配置.分析线段中枢 {
|
||||
中枢::分析(&self.线段序列组[i], &mut self.中枢序列组[i], true, "", 0);
|
||||
}
|
||||
}
|
||||
中枢::分析(&self.线段序列组[i], &mut self.中枢序列组[i], true, "", 0);
|
||||
}
|
||||
|
||||
// Step 5: 扩展线段分析 — 3 级递归
|
||||
for i in 0..self.扩展线段分析层次 {
|
||||
if i == 0 {
|
||||
线段::扩展分析(&self.笔序列, &mut self.扩展线段序列组[i], &self.配置);
|
||||
} else {
|
||||
let 源序列 = self.扩展线段序列组[i - 1].clone();
|
||||
线段::扩展分析(&源序列, &mut self.扩展线段序列组[i], &self.配置);
|
||||
if self.配置.分析扩展线段 || self.配置.分析线段中枢 {
|
||||
for i in 0..self.扩展线段分析层次 {
|
||||
if i == 0 {
|
||||
if self.配置.分析扩展线段 {
|
||||
线段::扩展分析(&self.笔序列, &mut self.扩展线段序列组[i], &self.配置);
|
||||
}
|
||||
} else {
|
||||
if self.配置.分析扩展线段 {
|
||||
let 源序列 = self.扩展线段序列组[i - 1].clone();
|
||||
线段::扩展分析(&源序列, &mut self.扩展线段序列组[i], &self.配置);
|
||||
}
|
||||
}
|
||||
if self.配置.分析线段中枢 {
|
||||
中枢::分析(
|
||||
&self.扩展线段序列组[i],
|
||||
&mut self.扩展中枢序列组[i],
|
||||
true,
|
||||
"",
|
||||
0,
|
||||
);
|
||||
}
|
||||
}
|
||||
中枢::分析(
|
||||
&self.扩展线段序列组[i],
|
||||
&mut self.扩展中枢序列组[i],
|
||||
true,
|
||||
"",
|
||||
0,
|
||||
);
|
||||
}
|
||||
|
||||
// Step 6: 混合扩展线段分析 — 3 级递归 (源 = 线段序列组[i])
|
||||
// NOTE: 当 线段分析层次=0 时 线段序列组 为空,用 min 避免越界
|
||||
for i in 0..self.混合扩展线段分析层次.min(self.线段序列组.len()) {
|
||||
let 源序列 = self.线段序列组[i].clone();
|
||||
线段::扩展分析(&源序列, &mut self.混合扩展线段序列组[i], &self.配置);
|
||||
中枢::分析(
|
||||
&self.混合扩展线段序列组[i],
|
||||
&mut self.混合扩展中枢序列组[i],
|
||||
true,
|
||||
"",
|
||||
0,
|
||||
);
|
||||
// Step 6: 混合扩展线段分析 — 3 级递归
|
||||
if self.配置.分析扩展线段 || self.配置.分析线段中枢 {
|
||||
for i in 0..self.混合扩展线段分析层次.min(self.线段序列组.len()) {
|
||||
if self.配置.分析扩展线段 {
|
||||
let 源序列 = self.线段序列组[i].clone();
|
||||
线段::扩展分析(&源序列, &mut self.混合扩展线段序列组[i], &self.配置);
|
||||
}
|
||||
if self.配置.分析线段中枢 {
|
||||
中枢::分析(
|
||||
&self.混合扩展线段序列组[i],
|
||||
&mut self.混合扩展中枢序列组[i],
|
||||
true,
|
||||
"",
|
||||
0,
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -366,73 +391,103 @@ impl 观察者 {
|
||||
self.混合扩展中枢序列组.push(Vec::new());
|
||||
}
|
||||
|
||||
for i in 1..self.缠论K线序列.len() - 1 {
|
||||
let 当前分型 = 分型::new(
|
||||
Some(Arc::clone(&self.缠论K线序列[i - 1])),
|
||||
Arc::clone(&self.缠论K线序列[i]),
|
||||
Some(Arc::clone(&self.缠论K线序列[i + 1])),
|
||||
);
|
||||
笔::分析(
|
||||
Arc::new(当前分型),
|
||||
&mut self.分型序列,
|
||||
&mut self.笔序列,
|
||||
&self.缠论K线序列,
|
||||
&self.普通K线序列,
|
||||
0,
|
||||
&self.配置,
|
||||
);
|
||||
}
|
||||
|
||||
中枢::分析(&self.笔序列, &mut self.笔_中枢序列, true, "", 0);
|
||||
|
||||
for i in 0..self.线段分析层次 {
|
||||
if i == 0 {
|
||||
线段::分析(
|
||||
&self.笔序列,
|
||||
&mut self.线段序列组[i],
|
||||
&self.配置,
|
||||
0,
|
||||
&[相对方向::向上, 相对方向::向下],
|
||||
if self.配置.分析笔 {
|
||||
for i in 1..self.缠论K线序列.len() - 1 {
|
||||
let 当前分型 = 分型::new(
|
||||
Some(Arc::clone(&self.缠论K线序列[i - 1])),
|
||||
Arc::clone(&self.缠论K线序列[i]),
|
||||
Some(Arc::clone(&self.缠论K线序列[i + 1])),
|
||||
);
|
||||
} else {
|
||||
let 源序列 = self.线段序列组[i - 1].clone();
|
||||
线段::分析(
|
||||
&源序列,
|
||||
&mut self.线段序列组[i],
|
||||
&self.配置,
|
||||
笔::分析(
|
||||
Arc::new(当前分型),
|
||||
&mut self.分型序列,
|
||||
&mut self.笔序列,
|
||||
&self.缠论K线序列,
|
||||
&self.普通K线序列,
|
||||
0,
|
||||
&[相对方向::向上, 相对方向::向下],
|
||||
&self.配置,
|
||||
);
|
||||
}
|
||||
中枢::分析(&self.线段序列组[i], &mut self.中枢序列组[i], true, "", 0);
|
||||
}
|
||||
|
||||
for i in 0..self.扩展线段分析层次 {
|
||||
if i == 0 {
|
||||
线段::扩展分析(&self.笔序列, &mut self.扩展线段序列组[i], &self.配置);
|
||||
} else {
|
||||
let 源序列 = self.扩展线段序列组[i - 1].clone();
|
||||
线段::扩展分析(&源序列, &mut self.扩展线段序列组[i], &self.配置);
|
||||
if self.笔序列.is_empty() {
|
||||
return;
|
||||
}
|
||||
|
||||
if self.配置.分析笔中枢 {
|
||||
中枢::分析(&self.笔序列, &mut self.笔_中枢序列, true, "", 0);
|
||||
}
|
||||
|
||||
if self.配置.分析线段 || self.配置.分析线段中枢 {
|
||||
for i in 0..self.线段分析层次 {
|
||||
if i == 0 {
|
||||
if self.配置.分析线段 {
|
||||
线段::分析(
|
||||
&self.笔序列,
|
||||
&mut self.线段序列组[i],
|
||||
&self.配置,
|
||||
0,
|
||||
&[相对方向::向上, 相对方向::向下],
|
||||
);
|
||||
}
|
||||
} else {
|
||||
if self.配置.分析线段 {
|
||||
let 源序列 = self.线段序列组[i - 1].clone();
|
||||
线段::分析(
|
||||
&源序列,
|
||||
&mut self.线段序列组[i],
|
||||
&self.配置,
|
||||
0,
|
||||
&[相对方向::向上, 相对方向::向下],
|
||||
);
|
||||
}
|
||||
}
|
||||
if self.配置.分析线段中枢 {
|
||||
中枢::分析(&self.线段序列组[i], &mut self.中枢序列组[i], true, "", 0);
|
||||
}
|
||||
}
|
||||
中枢::分析(
|
||||
&self.扩展线段序列组[i],
|
||||
&mut self.扩展中枢序列组[i],
|
||||
true,
|
||||
"",
|
||||
0,
|
||||
);
|
||||
}
|
||||
|
||||
for i in 0..self.混合扩展线段分析层次.min(self.线段序列组.len()) {
|
||||
let 源序列 = self.线段序列组[i].clone();
|
||||
线段::扩展分析(&源序列, &mut self.混合扩展线段序列组[i], &self.配置);
|
||||
中枢::分析(
|
||||
&self.混合扩展线段序列组[i],
|
||||
&mut self.混合扩展中枢序列组[i],
|
||||
true,
|
||||
"",
|
||||
0,
|
||||
);
|
||||
if self.配置.分析扩展线段 || self.配置.分析线段中枢 {
|
||||
for i in 0..self.扩展线段分析层次 {
|
||||
if i == 0 {
|
||||
if self.配置.分析扩展线段 {
|
||||
线段::扩展分析(&self.笔序列, &mut self.扩展线段序列组[i], &self.配置);
|
||||
}
|
||||
} else {
|
||||
if self.配置.分析扩展线段 {
|
||||
let 源序列 = self.扩展线段序列组[i - 1].clone();
|
||||
线段::扩展分析(&源序列, &mut self.扩展线段序列组[i], &self.配置);
|
||||
}
|
||||
}
|
||||
if self.配置.分析线段中枢 {
|
||||
中枢::分析(
|
||||
&self.扩展线段序列组[i],
|
||||
&mut self.扩展中枢序列组[i],
|
||||
true,
|
||||
"",
|
||||
0,
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if self.配置.分析扩展线段 || self.配置.分析线段中枢 {
|
||||
for i in 0..self.混合扩展线段分析层次.min(self.线段序列组.len()) {
|
||||
if self.配置.分析扩展线段 {
|
||||
let 源序列 = self.线段序列组[i].clone();
|
||||
线段::扩展分析(&源序列, &mut self.混合扩展线段序列组[i], &self.配置);
|
||||
}
|
||||
if self.配置.分析线段中枢 {
|
||||
中枢::分析(
|
||||
&self.混合扩展线段序列组[i],
|
||||
&mut self.混合扩展中枢序列组[i],
|
||||
true,
|
||||
"",
|
||||
0,
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -582,6 +637,97 @@ impl 观察者 {
|
||||
self.配置 = 配置;
|
||||
self.加载本地数据(文件路径)
|
||||
}
|
||||
|
||||
/// 相等 — 全量序列逐项比对,双端一致性验证,对应 Python `观察者相等`
|
||||
pub fn 相等(&self, other: &Self, 浮点容差: f64) -> (bool, String) {
|
||||
let 标签 = format!("观察者校验[A={},B={}]", self.标识(), other.标识());
|
||||
|
||||
if self.缠论K线序列.len() != other.缠论K线序列.len() {
|
||||
return (
|
||||
false,
|
||||
format!(
|
||||
"{标签}: 缠K序列长度不一致 A={},B={}",
|
||||
self.缠论K线序列.len(),
|
||||
other.缠论K线序列.len()
|
||||
),
|
||||
);
|
||||
}
|
||||
if self.分型序列.len() != other.分型序列.len() {
|
||||
return (
|
||||
false,
|
||||
format!(
|
||||
"{标签}: 分型序列长度不一致 A={},B={}",
|
||||
self.分型序列.len(),
|
||||
other.分型序列.len()
|
||||
),
|
||||
);
|
||||
}
|
||||
if self.笔序列.len() != other.笔序列.len() {
|
||||
return (
|
||||
false,
|
||||
format!(
|
||||
"{标签}: 笔序列长度不一致 A={},B={}",
|
||||
self.笔序列.len(),
|
||||
other.笔序列.len()
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
for (i, (a, b)) in self.笔序列.iter().zip(other.笔序列.iter()).enumerate() {
|
||||
let (eq, msg) = a.相等(b, 浮点容差);
|
||||
if !eq {
|
||||
return (false, format!("{标签}: 笔#{i}不一致 >> {msg}"));
|
||||
}
|
||||
}
|
||||
|
||||
if self.笔_中枢序列.len() != other.笔_中枢序列.len() {
|
||||
return (false, format!("{标签}: 笔中枢序列长度不一致"));
|
||||
}
|
||||
for (i, (a, b)) in self
|
||||
.笔_中枢序列
|
||||
.iter()
|
||||
.zip(other.笔_中枢序列.iter())
|
||||
.enumerate()
|
||||
{
|
||||
let (eq, msg) = a.相等(b, 浮点容差);
|
||||
if !eq {
|
||||
return (false, format!("{标签}: 笔中枢#{i}不一致 >> {msg}"));
|
||||
}
|
||||
}
|
||||
|
||||
for level in 0..self.线段分析层次.min(other.线段分析层次) {
|
||||
let a_segs = &self.线段序列组[level];
|
||||
let b_segs = &other.线段序列组[level];
|
||||
if a_segs.len() != b_segs.len() {
|
||||
return (false, format!("{标签}: 线段序列组[{level}]长度不一致"));
|
||||
}
|
||||
for (i, (a, b)) in a_segs.iter().zip(b_segs.iter()).enumerate() {
|
||||
let (eq, msg) = a.相等(b, 浮点容差);
|
||||
if !eq {
|
||||
return (
|
||||
false,
|
||||
format!("{标签}: 线段序列组[{level}]#{i}不一致 >> {msg}"),
|
||||
);
|
||||
}
|
||||
}
|
||||
let a_hubs = &self.中枢序列组[level];
|
||||
let b_hubs = &other.中枢序列组[level];
|
||||
if a_hubs.len() != b_hubs.len() {
|
||||
return (false, format!("{标签}: 中枢序列组[{level}]长度不一致"));
|
||||
}
|
||||
for (i, (a, b)) in a_hubs.iter().zip(b_hubs.iter()).enumerate() {
|
||||
let (eq, msg) = a.相等(b, 浮点容差);
|
||||
if !eq {
|
||||
return (
|
||||
false,
|
||||
format!("{标签}: 中枢序列组[{level}]#{i}不一致 >> {msg}"),
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
(true, format!("{标签}:全量序列校验全部一致"))
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
|
||||
@@ -31,11 +31,19 @@ pub struct K线合成器 {
|
||||
pub 周期组: Vec<i64>,
|
||||
pub 当前K线: HashMap<i64, Option<K线>>,
|
||||
pub 合成K线列表: HashMap<i64, Vec<K线>>,
|
||||
/// 事件回调 — K线完成时触发,对应 Python K线合成器.事件回调
|
||||
/// 签名: fn(信号类型: str, 标识: str, 周期: i64, 完成K线: K线)
|
||||
/// 在 _完成K线 清空当前K线后、新K线创建前触发
|
||||
事件回调: Option<Box<dyn Fn(String, String, i64, K线) + Send + Sync>>,
|
||||
}
|
||||
|
||||
impl K线合成器 {
|
||||
/// 创建K线合成器,按周期升序排列,初始化当前K线和合成K线列表
|
||||
pub fn new(标识: String, 周期组: Vec<i64>) -> Self {
|
||||
/// 创建K线合成器 — 对应 Python K线合成器.__init__(标识, 周期组, 事件回调=None)
|
||||
pub fn new(
|
||||
标识: String,
|
||||
周期组: Vec<i64>,
|
||||
事件回调: Option<Box<dyn Fn(String, String, i64, K线) + Send + Sync>>,
|
||||
) -> Self {
|
||||
let mut 周期组 = 周期组;
|
||||
周期组.sort();
|
||||
|
||||
@@ -51,37 +59,33 @@ impl K线合成器 {
|
||||
周期组,
|
||||
当前K线,
|
||||
合成K线列表,
|
||||
事件回调,
|
||||
}
|
||||
}
|
||||
|
||||
/// 设置事件回调 — 对应 Python `设置事件回调`
|
||||
pub fn 设置事件回调(
|
||||
&mut self,
|
||||
回调: Box<dyn Fn(String, String, i64, K线) + Send + Sync>,
|
||||
) {
|
||||
self.事件回调 = Some(回调);
|
||||
}
|
||||
|
||||
/// 投喂 — 便捷入口,直接从 OHLCV 创建 K线 并投喂
|
||||
pub fn 投喂(
|
||||
&mut self,
|
||||
时间戳: i64,
|
||||
开: f64,
|
||||
高: f64,
|
||||
低: f64,
|
||||
收: f64,
|
||||
量: f64,
|
||||
) -> Vec<(i64, K线)> {
|
||||
pub fn 投喂(&mut self, 时间戳: i64, 开: f64, 高: f64, 低: f64, 收: f64, 量: f64) {
|
||||
let 普K = K线::创建普K(&self.标识, 时间戳, 开, 高, 低, 收, 量, 0, 0);
|
||||
self.投喂K线(普K)
|
||||
self.投喂K线(普K);
|
||||
}
|
||||
|
||||
/// 投喂K线 — 输入最小周期K线,合成为所有目标周期
|
||||
/// 返回本次投喂完成了哪些周期的K线(周期 → 完成K线)
|
||||
pub fn 投喂K线(&mut self, 普K: K线) -> Vec<(i64, K线)> {
|
||||
let mut 完成记录 = Vec::new();
|
||||
pub fn 投喂K线(&mut self, 普K: K线) {
|
||||
let 周期组 = self.周期组.clone();
|
||||
for 周期 in 周期组 {
|
||||
if let Some(完成K线) = self._处理单个周期(周期, &普K) {
|
||||
完成记录.push((周期, 完成K线));
|
||||
}
|
||||
self._处理单个周期(周期, &普K);
|
||||
}
|
||||
完成记录
|
||||
}
|
||||
|
||||
fn _处理单个周期(&mut self, 周期: i64, 普K: &K线) -> Option<K线> {
|
||||
fn _处理单个周期(&mut self, 周期: i64, 普K: &K线) {
|
||||
let 目标时间戳 = self._对齐时间戳(普K.时间戳, 周期);
|
||||
let 相同时间 = self.当前K线[&周期]
|
||||
.as_ref()
|
||||
@@ -91,19 +95,17 @@ impl K线合成器 {
|
||||
if self.当前K线[&周期].is_none() {
|
||||
let 新K线 = self._创建新K线(周期, 目标时间戳, 普K);
|
||||
self.当前K线.insert(周期, Some(新K线));
|
||||
None
|
||||
} else if 相同时间 {
|
||||
let ent = self.当前K线.get_mut(&周期).unwrap();
|
||||
Self::_更新K线(ent.as_mut().unwrap(), 普K);
|
||||
None
|
||||
} else {
|
||||
let 完成K线 = self._完成K线(周期);
|
||||
self._完成K线(周期);
|
||||
let 新K线 = self._创建新K线(周期, 目标时间戳, 普K);
|
||||
self.当前K线.insert(周期, Some(新K线));
|
||||
完成K线
|
||||
}
|
||||
}
|
||||
|
||||
/// 对齐时间戳到周期边界 — 对应 Python `_对齐时间戳`
|
||||
fn _对齐时间戳(&self, 时间戳: i64, 周期: i64) -> i64 {
|
||||
if 周期 == 0 {
|
||||
panic!("_对齐时间戳: 周期不能为0");
|
||||
@@ -111,6 +113,7 @@ impl K线合成器 {
|
||||
(时间戳 / 周期) * 周期
|
||||
}
|
||||
|
||||
/// 创建新K线 — 对应 Python `_创建新K线`
|
||||
fn _创建新K线(&self, 周期: i64, 时间戳: i64, 普K: &K线) -> K线 {
|
||||
let 序号 = self
|
||||
.合成K线列表
|
||||
@@ -132,6 +135,7 @@ impl K线合成器 {
|
||||
)
|
||||
}
|
||||
|
||||
/// 更新K线 — 对应 Python `_更新K线`
|
||||
fn _更新K线(当前K线: &mut K线, 新数据: &K线) {
|
||||
当前K线.高 = 当前K线.高.max(新数据.高);
|
||||
当前K线.低 = 当前K线.低.min(新数据.低);
|
||||
@@ -139,9 +143,14 @@ impl K线合成器 {
|
||||
当前K线.成交量 += 新数据.成交量;
|
||||
}
|
||||
|
||||
fn _完成K线(&mut self, 周期: i64) -> Option<K线> {
|
||||
/// 完成K线 — 对应 Python `_完成K线`
|
||||
/// 清空当前K线后,触发事件回调(此时获取当前K线返回 None)
|
||||
fn _完成K线(&mut self, 周期: i64) {
|
||||
let ent = self.当前K线.get_mut(&周期).unwrap();
|
||||
let mut k线 = ent.take()?;
|
||||
let mut k线 = match ent.take() {
|
||||
Some(k) => k,
|
||||
None => return,
|
||||
};
|
||||
k线.序号 = self
|
||||
.合成K线列表
|
||||
.get(&周期)
|
||||
@@ -151,11 +160,130 @@ impl K线合成器 {
|
||||
|
||||
let 完成K线 = k线.clone();
|
||||
self.合成K线列表.get_mut(&周期).unwrap().push(k线);
|
||||
Some(完成K线)
|
||||
|
||||
// 对应 Python _完成K线:清空当前K线后、新K线创建前触发回调
|
||||
self._产生完成K线信号(周期, 完成K线);
|
||||
}
|
||||
|
||||
/// 获取指定周期当前正在合成的K线
|
||||
/// 产生完成K线信号 — 对应 Python `_产生完成K线信号`
|
||||
fn _产生完成K线信号(&self, 周期: i64, 完成K线: K线) {
|
||||
if let Some(ref cb) = self.事件回调 {
|
||||
cb("K线完成".into(), self.标识.clone(), 周期, 完成K线);
|
||||
}
|
||||
}
|
||||
|
||||
/// 获取指定周期当前正在合成的K线 — 对应 Python `获取当前K线`
|
||||
pub fn 获取当前K线(&self, 周期: i64) -> Option<&K线> {
|
||||
self.当前K线.get(&周期).and_then(|k| k.as_ref())
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn test_创建合成器_初始状态正确() {
|
||||
let synth = K线合成器::new("btcusd".into(), vec![60, 300], None);
|
||||
assert_eq!(synth.标识, "btcusd");
|
||||
assert_eq!(synth.周期组, vec![60, 300]);
|
||||
assert!(synth.事件回调.is_none());
|
||||
assert!(synth.当前K线[&60].is_none());
|
||||
assert!(synth.当前K线[&300].is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_设置事件回调() {
|
||||
let mut synth = K线合成器::new("btcusd".into(), vec![60], None);
|
||||
assert!(synth.事件回调.is_none());
|
||||
synth.设置事件回调(Box::new(|_, _, _, _| {}));
|
||||
assert!(synth.事件回调.is_some());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_对齐时间戳() {
|
||||
let synth = K线合成器::new("t".into(), vec![300], None);
|
||||
assert_eq!(synth._对齐时间戳(1218124800, 300), 1218124800);
|
||||
assert_eq!(synth._对齐时间戳(1218124801, 300), 1218124800);
|
||||
assert_eq!(synth._对齐时间戳(1218125099, 300), 1218124800);
|
||||
assert_eq!(synth._对齐时间戳(1218125100, 300), 1218125100);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_创建新K线_序号递进() {
|
||||
let mut synth = K线合成器::new("btcusd".into(), vec![300], None);
|
||||
{
|
||||
let first = K线::创建普K("btcusd", 0, 100.0, 110.0, 90.0, 105.0, 1000.0, 0, 300);
|
||||
synth.合成K线列表.get_mut(&300).unwrap().push(first);
|
||||
}
|
||||
let new_bar = K线::创建普K("btcusd", 100, 200.0, 210.0, 190.0, 205.0, 500.0, 0, 60);
|
||||
let created = synth._创建新K线(300, 300, &new_bar);
|
||||
assert_eq!(created.序号, 1);
|
||||
assert_eq!(created.时间戳, 300);
|
||||
assert_eq!(created.开盘价, 200.0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_更新K线_高低更新() {
|
||||
let mut current = K线::创建普K("t", 0, 100.0, 110.0, 90.0, 105.0, 100.0, 0, 300);
|
||||
let new_data = K线::创建普K("t", 0, 102.0, 115.0, 85.0, 108.0, 50.0, 0, 60);
|
||||
K线合成器::_更新K线(&mut current, &new_data);
|
||||
assert_eq!(current.高, 115.0);
|
||||
assert_eq!(current.低, 85.0);
|
||||
assert_eq!(current.收盘价, 108.0);
|
||||
assert_eq!(current.成交量, 150.0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_完成K线_返回完成K并将当前置空() {
|
||||
let mut synth = K线合成器::new("btcusd".into(), vec![300], None);
|
||||
let bar = K线::创建普K("btcusd", 300, 100.0, 110.0, 90.0, 105.0, 1000.0, 0, 300);
|
||||
synth.当前K线.insert(300, Some(bar));
|
||||
synth._完成K线(300);
|
||||
assert!(synth.当前K线[&300].is_none());
|
||||
assert_eq!(synth.合成K线列表[&300].len(), 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_完成K线_事件回调触发() {
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
|
||||
let callback_fired = Arc::new(AtomicBool::new(false));
|
||||
let cb_flag = Arc::clone(&callback_fired);
|
||||
|
||||
let mut synth = K线合成器::new(
|
||||
"btcusd".into(),
|
||||
vec![300],
|
||||
Some(Box::new(move |信号类型, 标识, 周期, _完成K线| {
|
||||
assert_eq!(信号类型, "K线完成");
|
||||
assert_eq!(标识, "btcusd");
|
||||
assert_eq!(周期, 300);
|
||||
cb_flag.store(true, Ordering::SeqCst);
|
||||
})),
|
||||
);
|
||||
|
||||
let bar1 = K线::创建普K("btcusd", 0, 100.0, 110.0, 90.0, 105.0, 1000.0, 0, 300);
|
||||
synth.当前K线.insert(300, Some(bar1));
|
||||
let bar2 = K线::创建普K("btcusd", 400, 200.0, 210.0, 190.0, 205.0, 500.0, 0, 60);
|
||||
synth.投喂K线(bar2);
|
||||
assert!(callback_fired.load(Ordering::SeqCst));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_投喂K线_多周期合成() {
|
||||
let mut synth = K线合成器::new("btcusd".into(), vec![60, 300], None);
|
||||
synth.投喂K线(K线::创建普K(
|
||||
"btcusd", 60, 100.0, 110.0, 90.0, 105.0, 100.0, 0, 60,
|
||||
));
|
||||
assert!(synth.获取当前K线(60).is_some());
|
||||
assert!(synth.获取当前K线(300).is_some());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_投喂_便捷方法() {
|
||||
let mut synth = K线合成器::new("btcusd".into(), vec![300], None);
|
||||
synth.投喂(1218124800, 100.0, 110.0, 90.0, 105.0, 1000.0);
|
||||
assert!(synth.获取当前K线(300).is_some());
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user