use base64::{engine::general_purpose, Engine};
use chrono::TimeZone;
use rsa::{
pkcs8::{DecodePrivateKey, DecodePublicKey},
Pkcs1v15Sign,
};
use std::time::{SystemTime, UNIX_EPOCH};
use reqwest::Client;
use serde::{Deserialize, Serialize};
use thiserror::Error;
#[derive(Error, Debug)]
pub enum WxpayApiError {
#[error("timestamp is invalid")]
InvalidTimestamp,
#[error("invalid ciphertext: {0}")]
InvalidCiphertext(#[from] aes_gcm::aes::cipher::InvalidLength),
#[error("invalid public key")]
InvalidPublicKey,
#[error("invalid private key")]
InvalidPrivateKey,
#[error("decrypt failed")]
DecryptFailed,
#[error("request error: {0}")]
RequestErr(#[from] reqwest::Error),
#[error("base64 decode error: {0}")]
Base64DecodeError(#[from] base64::DecodeError),
}
#[derive(Debug, Deserialize, Serialize)]
pub struct BatchTransferRequest {
pub appid: String,
pub out_batch_no: String,
pub batch_name: String,
pub batch_remark: String,
pub total_amount: u64,
pub total_num: u64,
pub transfer_detail_list: Vec<TransferDetail>,
pub transfer_scene_id: Option<String>,
pub notify_url: Option<String>,
}
#[derive(Debug, Deserialize, Serialize)]
pub struct TransferDetail {
pub out_detail_no: String,
pub transfer_amount: u64,
pub transfer_remark: String,
pub openid: String,
pub user_name: Option<String>,
}
#[derive(Debug, Deserialize, Serialize)]
pub struct BatchTransferResponse {
pub out_detail_no: String,
pub batch_id: String,
pub create_time: String,
pub out_batch_no: String,
}
#[derive(Debug, Deserialize, Serialize)]
pub struct WxPayFailedResponse {
pub code: String,
pub message: String,
pub detail: Option<serde_json::Value>,
}
#[derive(Debug, Deserialize, Serialize)]
#[serde(untagged)]
pub enum WxpayResponse<D> {
Failed(WxPayFailedResponse),
Success(D),
}
pub struct RequestBatchTransfer<'a> {
pub mchid: &'a str,
pub mch_private_key: &'a str,
pub mch_serial_no: &'a str,
pub wxpay_serial_no: &'a str,
pub request: BatchTransferRequest,
}
pub async fn request_batch_transfer<'a>(
RequestBatchTransfer {
mchid,
mch_private_key,
mch_serial_no,
wxpay_serial_no,
request,
}: RequestBatchTransfer<'a>,
) -> Result<WxpayResponse<serde_json::Value>, WxpayApiError> {
let url = "https://api.mch.weixin.qq.com/v3/transfer/batches";
let method = "POST";
let body = serde_json::to_string(&request).expect("failed to serialize request");
let (signature, timestamp, nonce_str) =
generate_wxpay_signature(method, "/v3/transfer/batches", mch_private_key, Some(&body))?;
let client = Client::new();
let response = client.post(url)
.header("Content-Type", "application/json")
.header("Accept", "application/json")
.header("User-Agent", "wechat-vendor-sdk/0.0.0")
.header("Wechatpay-Serial", wxpay_serial_no)
.header("Authorization", format!("WECHATPAY2-SHA256-RSA2048 mchid=\"{}\",nonce_str=\"{}\",signature=\"{}\",timestamp=\"{}\",serial_no=\"{}\"",
mchid, nonce_str, signature, timestamp, mch_serial_no))
.body(body)
.send()
.await?;
let result: WxpayResponse<serde_json::Value> = response.json().await?;
Ok(result)
}
pub fn generate_wxpay_signature(
method: &str,
url_path: &str,
private_key: &str,
body: Option<&str>,
) -> Result<(String, String, String), WxpayApiError> {
let timestamp = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_secs()
.to_string();
let nonce_str = generate_noncestr(32);
let mut content_to_sign = format!("{}\n{}\n{}\n{}\n", method, url_path, timestamp, nonce_str);
if let Some(body_content) = body {
content_to_sign.push_str(body_content);
content_to_sign.push('\n');
}
let signature = sha256_ras_and_base64(private_key, &content_to_sign)?;
Ok((signature, timestamp, nonce_str))
}
pub fn generate_noncestr(length: usize) -> String {
use rand::{distributions::Alphanumeric, Rng};
let noncestr: String = rand::thread_rng()
.sample_iter(&Alphanumeric)
.take(length)
.map(char::from)
.collect();
noncestr
}
pub fn sha256_ras_and_base64(
private_key: &str,
content_to_be_signed: &str,
) -> Result<String, WxpayApiError> {
use rsa::sha2::{Digest, Sha256};
use rsa::RsaPrivateKey;
let mut hasher = Sha256::new();
hasher.update(content_to_be_signed);
let hash = hasher.finalize();
let private_key =
RsaPrivateKey::from_pkcs8_pem(private_key).map_err(|e| WxpayApiError::InvalidPrivateKey)?;
let padding = Pkcs1v15Sign::new::<rsa::sha2::Sha256>();
let signature = private_key.sign(padding, &hash).expect("failed to sign");
let sig = general_purpose::STANDARD.encode(&signature);
Ok(sig)
}
pub fn verify_wxpay_callback_signature(
wx_public_key: &str,
signature: &str,
timestamp: &str,
nonce: &str,
body: &str,
skip_timestamp_check: Option<bool>,
) -> Result<bool, WxpayApiError> {
use base64::{engine::general_purpose, Engine as _};
use rsa::sha2::{Digest, Sha256};
use rsa::RsaPublicKey;
if skip_timestamp_check.unwrap_or(false) == false {
use chrono::{Duration, Utc};
let current_time = Utc::now();
let timestamp_secs = timestamp
.parse::<i64>()
.map_err(|_| WxpayApiError::InvalidTimestamp)?;
let timestamp_datetime = Utc
.timestamp_opt(timestamp_secs, 0)
.single()
.ok_or(WxpayApiError::InvalidTimestamp)?;
let time_diff = current_time - timestamp_datetime;
if time_diff > Duration::seconds(300) {
tracing::error!("timestamp is expired");
println!("timestamp is expired");
return Ok(false);
}
}
let signature_bytes = general_purpose::STANDARD.decode(signature)?;
let message = format!("{}\n{}\n{}\n", timestamp, nonce, body);
let mut hasher = Sha256::new();
hasher.update(message.as_bytes());
let hash = hasher.finalize();
let public_key = RsaPublicKey::from_public_key_pem(wx_public_key)
.map_err(|e| WxpayApiError::InvalidPublicKey)?;
let scheme = Pkcs1v15Sign::new::<Sha256>();
let res = public_key.verify(scheme, &hash, signature_bytes.as_slice());
match res {
Ok(_) => Ok(true),
Err(_) => Ok(false),
}
}
#[derive(Debug, Deserialize, Serialize)]
pub struct WxpayCallbackResourceDataClosed {
pub mchid: String,
pub out_batch_no: String,
pub batch_id: String,
pub batch_status: String,
pub total_num: i32,
pub total_amount: i32,
pub close_reason: Option<String>,
pub update_time: String,
}
#[derive(Debug, Deserialize, Serialize)]
pub struct WxpayCallbackResourceDataFinished {
pub out_batch_no: String,
pub batch_id: String,
pub batch_status: String,
pub total_amount: i32,
pub total_num: i32,
pub success_amount: i32,
pub success_num: i32,
pub fail_amount: i32,
pub fail_num: i32,
pub update_time: String,
}
#[derive(Debug, Deserialize, Serialize)]
#[serde(untagged)]
pub enum WxpayCallbackResourceData {
Closed(WxpayCallbackResourceDataClosed),
Finished(WxpayCallbackResourceDataFinished),
}
pub fn decrypt_wxpay_callback_resource(
apiv3_key: &str,
ciphertext: &str,
nonce: &str,
associated_data: &str,
) -> Result<WxpayCallbackResourceData, WxpayApiError> {
use aes_gcm::aead::{Aead, Payload};
use aes_gcm::{Aes256Gcm, KeyInit, Nonce};
use base64::{engine::general_purpose, Engine as _};
let ciphertext = general_purpose::STANDARD.decode(ciphertext)?;
let cipher = Aes256Gcm::new_from_slice(apiv3_key.as_bytes())?;
let payload = Payload {
msg: &ciphertext.as_slice(),
aad: &associated_data.as_bytes(),
};
let nonce = Nonce::from_slice(nonce.as_bytes());
let plaintext = cipher
.decrypt(nonce, payload)
.map_err(|_| WxpayApiError::DecryptFailed)?;
let val = serde_json::from_slice(&plaintext);
match val {
Ok(val) => Ok(val),
Err(e) => {
tracing::error!("decrypt wxpay callback resource to json failed: {}", e);
Err(WxpayApiError::DecryptFailed)
}
}
}
#[test]
fn test_generate_wxpay_signature() {
let private_key = "-----BEGIN PRIVATE KEY-----
...
-----END PRIVATE KEY-----
";
let method = "POST";
let url_path = "/v3/transfer/batches";
match generate_wxpay_signature(method, url_path, private_key, Some(r#"body"#)) {
Ok((signature, timestamp, nonce_str)) => {
println!("签名: {}", signature);
println!("时间戳: {}", timestamp);
println!("随机字符串: {}", nonce_str);
}
Err(e) => println!("生成签名失败: {}", e),
}
}
#[test]
fn test_verify_wxpay_callback_signature() {
let wx_public_key = "-----BEGIN PUBLIC KEY-----
...
-----END PUBLIC KEY-----";
let result = verify_wxpay_callback_signature(
wx_public_key,
"signature",
"1727611989",
"nonce",
r#"payload"#,
None,
);
println!("result: {:?}", result);
}
#[test]
fn test_timestamp_diff() {
use chrono::{Duration, Utc};
let current_time = Utc::now();
let timestamp_secs = ("1717233600")
.parse::<i64>()
.map_err(|e| format!("无效的时间戳: {}", e))
.unwrap();
let timestamp_datetime = Utc
.timestamp_opt(timestamp_secs, 0)
.single()
.ok_or("无效的时间戳")
.unwrap();
let time_diff = current_time - timestamp_datetime;
if time_diff > Duration::seconds(300) {
println!("duration11: {:?}", time_diff);
} else {
println!("duration22: {:?}", time_diff);
}
}
#[test]
fn test_decrypt_wxpay_callback_resource() {
let apiv3_key = "xxx";
let nonce = "yyy";
let ciphertext = "22qIR8j4SVcexi0PTqgsPPxXICxk+zz==";
let result = decrypt_wxpay_callback_resource(apiv3_key, ciphertext, nonce, "mch_payment");
println!("result: {:?}", result);
}