@@ -53,10 +53,14 @@ def _load_json(path: Path) -> dict[str, Any]:
def _load_required_env ( name : str , value : str | None ) - > str :
# 优先使用传入的 value,否则读取同名环境变量;两者均缺失时抛出 RuntimeError
resolved = value or os . environ . get ( name )
if not resolved :
raise RuntimeError ( f " Missing required value: { name } . " )
return resolved
if value :
return value
env_value = os . environ . get ( name )
if env_value :
return env_value
raise RuntimeError (
f " Missing required value ' { name } ' : not passed as argument and not set as environment variable. "
)
def default_output_dir ( ) - > Path :
@@ -67,6 +71,194 @@ def _maybe_path(enabled: bool, path: Path) -> Path | None:
return path if enabled else None
def _process_item (
* ,
index : int ,
item : Any ,
resolved_output_dir : Path ,
resolved_prompt_path : Path ,
resolved_run_id : str ,
debug_artifacts : bool ,
loaded_rules : list ,
filter_context : Any ,
max_retries : int ,
timeout_seconds : float ,
resolved_llm_api_key : str ,
resolved_llm_model : str ,
resolved_llm_api_url : str ,
) - > dict [ str , Any ] :
# 处理单条 item:提取 -> LLM 摘要 -> 规则过滤 -> 候选构建
# 返回 item_report dict; delivered 时额外携带 _candidate/_external_id 供调用方解包
# 步骤 1:确定各中间文件路径(debug_artifacts=False 时大部分路径为 None,不写盘)
item_key = f " item- { index : 02d } "
item_path = _maybe_path ( debug_artifacts , resolved_output_dir / " items " / f " { item_key } .item.json " )
extracted_path = resolved_output_dir / " extracted " / f " { item_key } .extracted.json "
summary_output = _maybe_path ( debug_artifacts , resolved_output_dir / " summary " / item_key / " result.loop.json " )
filter_path = _maybe_path ( debug_artifacts , resolved_output_dir / " filter " / f " { item_key } .filter.json " )
record_path = _maybe_path ( debug_artifacts , resolved_output_dir / " candidates " / f " { item_key } .article-candidate-record.json " )
openclaw_path = _maybe_path ( debug_artifacts , resolved_output_dir / " candidates " / f " { item_key } .openclaw-candidate-input.json " )
if item_path is not None :
_save_json ( item_path , item . model_dump ( mode = " json " ) )
# 步骤 2:初始化 item_report,记录基础元信息;debug 模式下附加各中间文件路径
item_report : dict [ str , Any ] = {
" item_key " : item_key ,
" item_id " : item . item_id ,
" external_id " : item . external_id ,
" url " : str ( item . url ) ,
" title " : item . title ,
" status " : " pulled " ,
}
if debug_artifacts :
item_report [ " paths " ] = {
" item " : str ( item_path ) if item_path else None ,
" extracted " : str ( extracted_path ) if extracted_path else None ,
" summary " : str ( summary_output ) if summary_output else None ,
" filter " : str ( filter_path ) if filter_path else None ,
" article_candidate " : str ( record_path ) if record_path else None ,
" openclaw_candidate " : str ( openclaw_path ) if openclaw_path else None ,
}
# 步骤 3:内容提取(RSS 内联内容 or 回源抓取);extracted.json 始终写盘
extraction = extract_content ( ExtractionInput ( item = item ) )
extracted_payload = extraction . model_dump ( mode = " json " )
if extracted_path is not None :
_save_json ( extracted_path , extracted_payload )
if not extraction . success or extraction . article is None :
item_report [ " status " ] = " extract_failed "
item_report [ " error " ] = extraction . error . model_dump ( mode = " json " ) if extraction . error else None
return item_report
# 步骤 4:LLM 摘要循环,失败时最多重试 max_retries 次
item_report [ " status " ] = " extracted "
summary_exit_code , summary_payload , summary_report = run_loop_payload (
extracted_payload = extracted_payload ,
prompt_path = resolved_prompt_path ,
output_path = summary_output ,
max_retries = max_retries ,
timeout_seconds = timeout_seconds ,
api_key = resolved_llm_api_key ,
model = resolved_llm_model ,
api_url = resolved_llm_api_url ,
)
if summary_exit_code != 0 or summary_payload is None :
item_report [ " status " ] = " summary_failed "
if summary_report is not None :
item_report [ " summary_errors " ] = summary_report . errors
return item_report
# 步骤 5:规则引擎过滤,产出 keep/review/drop 决策及 digest_rank
summary = LlmSummaryResult . model_validate ( summary_payload )
decision = evaluate_filter_rules (
FilterInput ( item = item , article = extraction . article , summary = summary , context = filter_context ) ,
loaded_rules ,
)
if filter_path is not None :
_save_json ( filter_path , decision . model_dump ( mode = " json " ) )
# 步骤 6:构建 ArticleCandidateRecord 和 OpenClawCandidateInput
record = build_article_candidate_record (
summary = summary ,
article = extraction . article ,
filter_result = decision ,
item = item ,
source_refs = CandidateSourceRefs (
item_path = str ( item_path ) if item_path else None ,
extracted_path = str ( extracted_path ) if extracted_path else None ,
summary_path = str ( summary_output ) if summary_output else None ,
filter_path = str ( filter_path ) if filter_path else None ,
) ,
metadata = CandidateMetadata (
generated_at = datetime . now ( tz = UTC ) ,
producer = " run_freshrss_pipeline " ,
run_id = resolved_run_id ,
) ,
)
openclaw_input = build_openclaw_candidate_input ( record )
if record_path is not None :
_save_json ( record_path , record . model_dump ( mode = " json " ) )
if openclaw_path is not None :
_save_json ( openclaw_path , openclaw_input . model_dump ( mode = " json " ) )
# 步骤 7:标记 delivered,附加临时键 _candidate/_external_id 供主函数解包后加入投递列表
item_report [ " status " ] = " delivered "
item_report [ " selection_decision " ] = decision . decision
item_report [ " candidate_id " ] = openclaw_input . candidate_id
item_report [ " _candidate " ] = openclaw_input
item_report [ " _external_id " ] = item . external_id
return item_report
def _build_and_persist_delivery (
* ,
delivered_candidates : list [ OpenClawCandidateInput ] ,
resolved_run_id : str ,
resolved_delivery_date : Any ,
delivery_output : Path ,
) - > tuple [ OpenClawDeliveryPayload , dict [ str , Any ] ] :
# 按 digest_rank 排序,构建 delivery payload,持久化词元索引
delivered_candidates . sort ( key = lambda candidate : candidate . digest_rank , reverse = True )
delivery_payload = build_openclaw_delivery_payload (
delivered_candidates ,
run_id = resolved_run_id ,
for_date = resolved_delivery_date ,
)
_save_json ( delivery_output , delivery_payload . model_dump ( mode = " json " ) )
keyword_index_result = persist_keyword_indexes (
delivery_payload . candidates ,
for_date = delivery_payload . date ,
digest_id = delivery_payload . run_id ,
source = " openclaw_delivery_payload " ,
daily_dir = DEFAULT_TERM_DAILY_DIR ,
stats_path = DEFAULT_TERM_STATS_PATH ,
aliases_path = DEFAULT_TERM_ALIASES_PATH ,
stopwords_path = DEFAULT_TERM_STOPWORDS_PATH ,
)
return delivery_payload , keyword_index_result
def _build_run_report (
* ,
resolved_run_id : str ,
started_at : datetime ,
limit : int ,
items : list ,
delivered_candidates : list ,
marked_count : int ,
mark_read : bool ,
debug_artifacts : bool ,
raw_output : Path ,
delivery_output : Path ,
keyword_index_result : dict [ str , Any ] ,
item_reports : list [ dict [ str , Any ] ] ,
) - > dict [ str , Any ] :
# 统计各状态计数,组装 run report dict
status_counts : dict [ str , int ] = { }
for item_report in item_reports :
status = str ( item_report [ " status " ] )
status_counts [ status ] = status_counts . get ( status , 0 ) + 1
return {
" run_id " : resolved_run_id ,
" started_at " : started_at . isoformat ( ) ,
" completed_at " : datetime . now ( tz = UTC ) . isoformat ( ) ,
" requested_limit " : limit ,
" pulled_count " : len ( items ) ,
" delivered_count " : len ( delivered_candidates ) ,
" marked_read_count " : marked_count ,
" mark_read_requested " : mark_read ,
" debug_artifacts " : debug_artifacts ,
" raw_output " : str ( raw_output ) ,
" delivery_output " : str ( delivery_output ) ,
" keyword_index " : keyword_index_result ,
" status_counts " : status_counts ,
" items " : item_reports ,
}
def run_freshrss_pipeline (
* ,
api_base_url : str | None = None ,
@@ -80,6 +272,7 @@ def run_freshrss_pipeline(
debug_artifacts : bool = False ,
prompt : Path | None = None ,
rules : Path | None = None ,
context : dict [ str , Any ] | None = None ,
context_path : Path | None = None ,
max_retries : int = 2 ,
timeout_seconds : float = 60.0 ,
@@ -90,6 +283,7 @@ def run_freshrss_pipeline(
delivery_date : date | None = None ,
output_dir : Path | None = None ,
) - > dict [ str , Any ] :
# --- 阶段 1:初始化 run_id、输出路径、凭证 ---
started_at = datetime . now ( tz = UTC )
resolved_output_dir = output_dir or default_output_dir ( )
run_stamp = started_at . strftime ( " % Y % m %d - % H % M % S " )
@@ -112,6 +306,7 @@ def run_freshrss_pipeline(
report_output = resolved_output_dir / " run-report.json "
items_list_output = _maybe_path ( debug_artifacts , resolved_output_dir / " items " / " freshrss.items.json " )
# --- 阶段 2:登录 FreshRSS,拉取未读条目,写原始 payload ---
client = FreshRSSClient (
api_base_url = resolved_api_base_url ,
username = resolved_username ,
@@ -137,158 +332,71 @@ def run_freshrss_pipeline(
_save_json ( items_list_output , [ item . model_dump ( mode = " json " ) for item in items ] )
loaded_rules = load_filter_rules ( resolved_rules_path )
context = FilterContext . model_validate ( _load_json ( context_path ) ) if context_path else FilterContext ( )
# context 优先使用直接传入的 dict,其次读取 context_path 文件,两者均缺失则使用空 context
if context is not None :
filter_context = FilterContext . model_validate ( context )
elif context_path is not None :
filter_context = FilterContext . model_validate ( _load_json ( context_path ) )
else :
filter_context = FilterContext ( )
# --- 阶段 3:逐条处理(提取 -> LLM 摘要 -> 规则过滤 -> 候选构建) ---
delivered_candidates : list [ OpenClawCandidateInput ] = [ ]
delivered_item_ids : list [ str ] = [ ]
item_reports : list [ dict [ str , Any ] ] = [ ]
for index , item in enumerate ( items , start = 1 ) :
# 逐条处理:提取 -> 摘要 -> 过滤 -> 候选构建;任一步骤失败则记录状态后 continue
item_key = f " item- { index : 02d } "
item_path = _maybe_path ( debug_artifacts , resolved_output_dir / " items " / f " { item_key } .item.json " )
extracted_path = resolved_output_dir / " extracted " / f " { item_key } .extracted.json "
summary_output = _maybe_path ( debug_artifacts , resolved_output_dir / " summary " / item_key / " result.loop.json " )
filter_path = _maybe_path ( debug_artifacts , resolved_output_dir / " filter " / f " { item_key } .filter.json " )
record_path = _maybe_path ( debug_artifacts , resolved_output_dir / " candidates " / f " { item_key } .article-candidate-record.json " )
openclaw_path = _maybe_path ( debug_artifacts , resolved_output_dir / " candidates " / f " { item_key } .openclaw-candidate-input.json " )
if item_path is not None :
_save_json ( item_path , item . model_dump ( mode = " json " ) )
item_report : dict [ str , Any ] = {
" item_key " : item_key ,
" item_id " : item . item_id ,
" external_id " : item . external_id ,
" url " : str ( item . url ) ,
" title " : item . title ,
" status " : " pulled " ,
}
if debug_artifacts :
item_report [ " paths " ] = {
" item " : str ( item_path ) if item_path else None ,
" extracted " : str ( extracted_path ) if extracted_path else None ,
" summary " : str ( summary_output ) if summary_output else None ,
" filter " : str ( filter_path ) if filter_path else None ,
" article_candidate " : str ( record_path ) if record_path else None ,
" openclaw_candidate " : str ( openclaw_path ) if openclaw_path else None ,
}
extraction = extract_content ( ExtractionInput ( item = item ) )
extracted_payload = extraction . model_dump ( mode = " json " )
if extracted_path is not None :
_save_json ( extracted_path , extracted_payload )
if not extraction . success or extraction . article is None :
item_report [ " status " ] = " extract_failed "
item_report [ " error " ] = extraction . error . model_dump ( mode = " json " ) if extraction . error else None
item_reports . append ( item_report )
continue
item_report [ " status " ] = " extracted "
summary_exit_code , summary_payload , summary_report = run_loop_payload (
extracted_payload = extracted_payload ,
prompt_path = resolved_prompt_path ,
output_path = summary_output ,
result = _process_item (
index = index ,
item = item ,
resolved_output_dir = resolved_output_dir ,
resolved_prompt_path = resolved_prompt_path ,
resolved_run_id = resolved_run_id ,
debug_artifacts = debug_artifacts ,
loaded_rules = loaded_rules ,
filter_context = filter_context ,
max_retries = max_retries ,
timeout_seconds = timeout_seconds ,
api_key = resolved_llm_api_key ,
model = resolved_llm_model ,
api_url = resolved_llm_api_url ,
resolved_llm_api_key = resolved_llm_api_key ,
resolved_llm_model = resolved_llm_model ,
resolved_llm_api_url = resolved_llm_api_url ,
)
if summary_exit_code ! = 0 or summary_payload is None :
item_report [ " status " ] = " summary_failed "
if summary_report is not None :
item_report [ " summary_errors " ] = summary_report . errors
item_reports . append ( item_report )
continue
if result . get ( " status " ) = = " delivered " :
delivered_candidates . append ( result . pop ( " _candidate " ) )
external_id = result . pop ( " _external_id " , None )
if external_id :
delivered_item_ids . append ( external_id )
item_reports . append ( result )
summary = LlmSummaryResult . model_validate ( summary_payload )
decision = evaluate_filter_rules (
FilterInput ( item = item , article = extraction . article , summary = summary , context = context ) ,
loaded_rules ,
)
if filter_path is not None :
_save_json ( filter_path , decision . model_dump ( mode = " json " ) )
record = build_article_candidate_record (
summary = summary ,
article = extraction . article ,
filter_result = decision ,
item = item ,
source_refs = CandidateSourceRefs (
item_path = str ( item_path ) if item_path else None ,
extracted_path = str ( extracted_path ) if extracted_path else None ,
summary_path = str ( summary_output ) if summary_output else None ,
filter_path = str ( filter_path ) if filter_path else None ,
) ,
metadata = CandidateMetadata (
generated_at = datetime . now ( tz = UTC ) ,
producer = " run_freshrss_pipeline " ,
run_id = resolved_run_id ,
) ,
)
openclaw_input = build_openclaw_candidate_input ( record )
if record_path is not None :
_save_json ( record_path , record . model_dump ( mode = " json " ) )
if openclaw_path is not None :
_save_json ( openclaw_path , openclaw_input . model_dump ( mode = " json " ) )
delivered_candidates . append ( openclaw_input )
if item . external_id :
delivered_item_ids . append ( item . external_id )
item_report [ " status " ] = " delivered "
item_report [ " selection_decision " ] = decision . decision
item_report [ " candidate_id " ] = openclaw_input . candidate_id
item_reports . append ( item_report )
delivered_candidates . sort ( key = lambda candidate : candidate . digest_rank , reverse = True )
delivery_payload = build_openclaw_delivery_payload (
delivered_candidates ,
run_id = resolved_run_id ,
for_date = resolved_delivery_date ,
)
_save_json ( delivery_output , delivery_payload . model_dump ( mode = " json " ) )
keyword_index_result = persist_keyword_indexes (
delivery_payload . candidates ,
for_date = delivery_payload . date ,
digest_id = delivery_payload . run_id ,
source = " openclaw_delivery_payload " ,
daily_dir = DEFAULT_TERM_DAILY_DIR ,
stats_path = DEFAULT_TERM_STATS_PATH ,
aliases_path = DEFAULT_TERM_ALIASES_PATH ,
stopwords_path = DEFAULT_TERM_STOPWORDS_PATH ,
# --- 阶段 4+5:构建 delivery payload 并持久化词元索引 ---
delivery_payload , keyword_index_result = _build_and_persist_delivery (
delivered_candidates = delivered_candidates ,
resolved_run_id = resolved_run_id ,
resolved_delivery_date = resolved_delivery_date ,
delivery_output = delivery_output ,
)
# --- 阶段 6:标记已读,汇总报告,返回结果 ---
marked_count = 0
if mark_read and delivered_item_ids :
# 仅标记成功投递(delivered)的 item; drop/review 的 item 保持未读状态
client . mark_items_as_read ( auth_token = auth_token , item_ids = delivered_item_ids )
marked_count = len ( { item_id for item_id in delivered_item_ids if item_id } )
status_counts : dict [ str , int ] = { }
for item_report in item_reports :
status = str ( item_report [ " status " ] )
status_counts [ status ] = status_counts . get ( status , 0 ) + 1
report = {
" run_id " : resolved_run_id ,
" started_at " : started_at . isoformat ( ) ,
" completed_at " : datetime . now ( tz = UTC ) . isoformat ( ) ,
" requested_limit " : limit ,
" pulled_count " : len ( items ) ,
" delivered_count " : len ( delivered_candidates ) ,
" marked_read_count " : marked_count ,
" mark_read_requested " : mark_read ,
" debug_artifacts " : debug_artifacts ,
" raw_output " : str ( raw_output ) ,
" delivery_output " : str ( delivery_output ) ,
" keyword_index " : keyword_index_result ,
" status_counts " : status_counts ,
" items " : item_reports ,
}
report = _build_run_report (
resolved_run_id = resolved_run_id ,
started_at = started_at ,
limit = limit ,
items = items ,
delivered_candidates = delivered_candidates ,
marked_count = marked_count ,
mark_read = mark_read ,
debug_artifacts = debug_artifacts ,
raw_output = raw_output ,
delivery_output = delivery_output ,
keyword_index_result = keyword_index_result ,
item_reports = item_reports ,
)
_save_json ( report_output , report )
return {
@@ -301,7 +409,7 @@ def run_freshrss_pipeline(
" pulled_count " : len ( items ) ,
" delivered_count " : len ( delivered_candidates ) ,
" marked_read_count " : marked_count ,
" status_counts " : status_counts ,
" status_counts " : report [ " status_counts " ] ,
" debug_artifacts " : debug_artifacts ,
" delivery_payload " : delivery_payload . model_dump ( mode = " json " ) ,
" items " : item_reports ,