Add transport raid functionality to Falukant daemon: Introduced falukant_transport_raid module with SQL queries for managing transport raids, including candidate transports and user reputation updates. Updated documentation to reflect new WebSocket events and integrated the module into the UndergroundWorker for seamless processing of transport-related events.
This commit is contained in:
505
src/worker/falukant_transport_raid.rs
Normal file
505
src/worker/falukant_transport_raid.rs
Normal file
@@ -0,0 +1,505 @@
|
||||
//! Transportüberfälle (`raid_transport`) und Wachen (`transport.guard_count`).
|
||||
//! Siehe `docs/FALUKANT_TRANSPORT_RAID_DAEMON.md`.
|
||||
//!
|
||||
//! Läuft im Underground-Tick (ca. 60 s). Kein 1-Tage-Wartezeit wie andere Untergrundjobs.
|
||||
|
||||
use crate::db::{ConnectionPool, DbError, Row};
|
||||
use crate::message_broker::MessageBroker;
|
||||
use rand::Rng;
|
||||
use serde_json::{json, Value as Json};
|
||||
|
||||
use crate::worker::sql::{
|
||||
QUERY_GET_AVAILABLE_STOCKS, QUERY_INSERT_INVENTORY, QUERY_RAID_ACTIVE_UNDERGROUND,
|
||||
QUERY_RAID_APP_USER_FOR_FALUKANT, QUERY_RAID_CANDIDATE_TRANSPORTS,
|
||||
QUERY_RAID_FALUKANT_USER_FOR_CHARACTER, QUERY_RAID_NEAREST_BRANCH_FOR_USER,
|
||||
QUERY_RAID_REGION_ALLOWED, QUERY_RAID_SCHEMA_READY, QUERY_RAID_SUBTRACT_REP_BY_USER,
|
||||
QUERY_RAID_UPDATE_UNDERGROUND_RESULT, QUERY_UPDATE_TRANSPORT_SIZE,
|
||||
};
|
||||
|
||||
/// Wird aus `UndergroundWorker::tick` aufgerufen.
|
||||
pub fn run_tick(pool: &ConnectionPool, broker: &MessageBroker) -> Result<(), DbError> {
|
||||
if !schema_ready(pool) {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let mut conn = pool
|
||||
.get()
|
||||
.map_err(|e| DbError::new(format!("DB-Verbindung fehlgeschlagen: {e}")))?;
|
||||
|
||||
conn.prepare("raid_active", QUERY_RAID_ACTIVE_UNDERGROUND)
|
||||
.map_err(|e| DbError::new(format!("prepare raid_active: {e}")))?;
|
||||
|
||||
let rows = conn
|
||||
.execute("raid_active", &[])
|
||||
.map_err(|e| DbError::new(format!("exec raid_active: {e}")))?;
|
||||
|
||||
drop(conn);
|
||||
|
||||
for row in rows {
|
||||
let ug_id = parse_i32(&row, "id", -1);
|
||||
let performer_cid = parse_i32(&row, "performer_id", -1);
|
||||
let params_str = row.get("parameters").cloned().unwrap_or_else(|| "{}".into());
|
||||
if ug_id < 0 || performer_cid < 0 {
|
||||
continue;
|
||||
}
|
||||
|
||||
if let Err(e) = try_resolve_one_raid(
|
||||
pool,
|
||||
broker,
|
||||
ug_id,
|
||||
performer_cid,
|
||||
¶ms_str,
|
||||
) {
|
||||
eprintln!("[falukant_transport_raid] underground_id={ug_id}: {e}");
|
||||
}
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn schema_ready(pool: &ConnectionPool) -> bool {
|
||||
let Ok(mut conn) = pool.get() else {
|
||||
return false;
|
||||
};
|
||||
if conn.prepare("rs", QUERY_RAID_SCHEMA_READY).is_err() {
|
||||
return false;
|
||||
}
|
||||
let Ok(rows) = conn.execute("rs", &[]) else {
|
||||
return false;
|
||||
};
|
||||
rows.first()
|
||||
.and_then(|r| r.get("ready"))
|
||||
.map(|v| v == "t" || v == "true")
|
||||
.unwrap_or(false)
|
||||
}
|
||||
|
||||
fn try_resolve_one_raid(
|
||||
pool: &ConnectionPool,
|
||||
broker: &MessageBroker,
|
||||
underground_id: i32,
|
||||
performer_character_id: i32,
|
||||
params_json: &str,
|
||||
) -> Result<(), DbError> {
|
||||
let p: Json = serde_json::from_str(params_json).unwrap_or_else(|_| json!({}));
|
||||
let region_id = p
|
||||
.get("regionId")
|
||||
.or_else(|| p.get("region_id"))
|
||||
.and_then(|v| v.as_i64())
|
||||
.map(|x| x as i32)
|
||||
.unwrap_or(-1);
|
||||
let band_size = p
|
||||
.get("bandSize")
|
||||
.or_else(|| p.get("band_size"))
|
||||
.and_then(|v| v.as_i64())
|
||||
.map(|x| x as i32)
|
||||
.unwrap_or(3)
|
||||
.clamp(1, 50);
|
||||
|
||||
if region_id < 0 {
|
||||
return finish_invalid(
|
||||
pool,
|
||||
broker,
|
||||
underground_id,
|
||||
performer_character_id,
|
||||
"invalid_parameters",
|
||||
json!({"message": "regionId required"}),
|
||||
);
|
||||
}
|
||||
|
||||
let mut conn = pool
|
||||
.get()
|
||||
.map_err(|e| DbError::new(format!("DB-Verbindung fehlgeschlagen: {e}")))?;
|
||||
conn.prepare("raid_reg", QUERY_RAID_REGION_ALLOWED)
|
||||
.map_err(|e| DbError::new(format!("prepare raid_reg: {e}")))?;
|
||||
let ok = conn
|
||||
.execute("raid_reg", &[®ion_id])
|
||||
.map_err(|e| DbError::new(format!("exec raid_reg: {e}")))?;
|
||||
if ok.is_empty() {
|
||||
return finish_invalid(
|
||||
pool,
|
||||
broker,
|
||||
underground_id,
|
||||
performer_character_id,
|
||||
"invalid_region",
|
||||
json!({"message": "region not allowed for raid_transport"}),
|
||||
);
|
||||
}
|
||||
|
||||
conn.prepare("raid_fu", QUERY_RAID_FALUKANT_USER_FOR_CHARACTER)
|
||||
.map_err(|e| DbError::new(format!("prepare raid_fu: {e}")))?;
|
||||
let fu_rows = conn
|
||||
.execute("raid_fu", &[&performer_character_id])
|
||||
.map_err(|e| DbError::new(format!("exec raid_fu: {e}")))?;
|
||||
let (performer_falukant_id, performer_app_uid) = fu_rows
|
||||
.first()
|
||||
.map(|r| {
|
||||
(
|
||||
parse_i32(r, "falukant_user_id", -1),
|
||||
parse_i32(r, "app_user_id", -1),
|
||||
)
|
||||
})
|
||||
.unwrap_or((-1, -1));
|
||||
if performer_falukant_id < 0 {
|
||||
return finish_invalid(
|
||||
pool,
|
||||
broker,
|
||||
underground_id,
|
||||
performer_character_id,
|
||||
"performer_not_found",
|
||||
json!({}),
|
||||
);
|
||||
}
|
||||
|
||||
drop(conn);
|
||||
|
||||
conn = pool
|
||||
.get()
|
||||
.map_err(|e| DbError::new(format!("DB-Verbindung fehlgeschlagen: {e}")))?;
|
||||
conn.prepare("raid_cand", QUERY_RAID_CANDIDATE_TRANSPORTS)
|
||||
.map_err(|e| DbError::new(format!("prepare raid_cand: {e}")))?;
|
||||
let candidates = conn
|
||||
.execute("raid_cand", &[®ion_id, &performer_falukant_id])
|
||||
.map_err(|e| DbError::new(format!("exec raid_cand: {e}")))?;
|
||||
|
||||
if candidates.is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let cand = &candidates[0];
|
||||
let transport_id = parse_i32(cand, "transport_id", -1);
|
||||
let transport_size = parse_i32(cand, "transport_size", 0);
|
||||
let product_id = parse_i32(cand, "product_id", -1);
|
||||
let guard_count = parse_i32(cand, "guard_count", 0).max(0);
|
||||
let victim_fu = parse_i32(cand, "victim_falukant_user_id", -1);
|
||||
let vehicle_cap = parse_i32(cand, "vehicle_capacity", 1).max(1);
|
||||
|
||||
if transport_id < 0 || product_id < 0 || victim_fu < 0 || transport_size < 2 {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let mut rng = rand::thread_rng();
|
||||
|
||||
// Begegnung (Bande + Sichtbarkeit großer Ladung)
|
||||
let visibility_penalty = ((transport_size as f64 / vehicle_cap as f64).min(2.0) - 1.0) * 0.04;
|
||||
let encounter_p = (0.08_f64 + band_size as f64 * 0.012 + visibility_penalty).clamp(0.05, 0.45);
|
||||
if !rng.gen_bool(encounter_p) {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
// Kampf
|
||||
let raid_power = band_size + rng.gen_range(0..=band_size);
|
||||
let gmax = guard_count.max(1);
|
||||
let guard_power = guard_count + rng.gen_range(0..=gmax);
|
||||
|
||||
let outcome = if raid_power > (guard_power as f64 * 1.25) as i32 {
|
||||
"major_success"
|
||||
} else if raid_power > guard_power {
|
||||
"partial_success"
|
||||
} else {
|
||||
"repelled"
|
||||
};
|
||||
|
||||
if outcome == "repelled" {
|
||||
let result = json!({
|
||||
"status": "completed",
|
||||
"bandSize": band_size,
|
||||
"regionId": region_id,
|
||||
"attempts": 1,
|
||||
"successes": 0,
|
||||
"lastTargetTransportId": transport_id,
|
||||
"lastLoot": Json::Null,
|
||||
"lastOutcome": "repelled"
|
||||
});
|
||||
save_result(pool, underground_id, &result)?;
|
||||
publish_underground(broker, performer_app_uid, "raid_repelled");
|
||||
publish_status(broker, performer_app_uid);
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
// Beuteanteil (nie Totalverlust)
|
||||
let base_loot = match outcome {
|
||||
"major_success" => rng.gen_range(0.35..0.60),
|
||||
_ => rng.gen_range(0.15..0.45),
|
||||
};
|
||||
let guard_mul = (1.0_f64 - 0.05_f64 * (guard_count.min(12) as f64)).max(0.35);
|
||||
let loot_share = (base_loot * guard_mul).clamp(0.05, 0.55);
|
||||
let mut stolen = ((transport_size as f64) * loot_share).floor() as i32;
|
||||
stolen = stolen.clamp(1, transport_size - 1);
|
||||
|
||||
let carry_loss = rng.gen_range(0.65..0.90);
|
||||
let mut to_store = (stolen as f64 * carry_loss).floor() as i32;
|
||||
to_store = to_store.max(0).min(stolen);
|
||||
|
||||
// Nächstgelegenes Lager
|
||||
let (branch_id, _) = nearest_branch(pool, region_id, performer_falukant_id)?;
|
||||
let (stored, lost_storage) = if branch_id > 0 && to_store > 0 {
|
||||
let remaining = add_to_branch_inventory(pool, branch_id, product_id, to_store)?;
|
||||
(to_store - remaining, remaining)
|
||||
} else {
|
||||
(0, to_store)
|
||||
};
|
||||
|
||||
let new_size = transport_size - stolen;
|
||||
if new_size < 1 {
|
||||
return Err(DbError::new("raid: invalid new transport size"));
|
||||
}
|
||||
|
||||
let mut conn = pool
|
||||
.get()
|
||||
.map_err(|e| DbError::new(format!("DB-Verbindung fehlgeschlagen: {e}")))?;
|
||||
conn.prepare("upd_sz", QUERY_UPDATE_TRANSPORT_SIZE)
|
||||
.map_err(|e| DbError::new(format!("prepare upd_sz: {e}")))?;
|
||||
conn.execute("upd_sz", &[&transport_id, &new_size])
|
||||
.map_err(|e| DbError::new(format!("exec upd_sz: {e}")))?;
|
||||
|
||||
// Opfer: kleiner Rufabzug
|
||||
conn.prepare("raid_rep", QUERY_RAID_SUBTRACT_REP_BY_USER)
|
||||
.map_err(|e| DbError::new(format!("prepare raid_rep: {e}")))?;
|
||||
let _ = conn.execute("raid_rep", &[&victim_fu, &2_i32]);
|
||||
|
||||
drop(conn);
|
||||
|
||||
let victim_app_uid = app_user_for_falukant(pool, victim_fu)?;
|
||||
|
||||
let result = json!({
|
||||
"status": "completed",
|
||||
"bandSize": band_size,
|
||||
"regionId": region_id,
|
||||
"attempts": 1,
|
||||
"successes": 1,
|
||||
"lastTargetTransportId": transport_id,
|
||||
"lastLoot": {
|
||||
"productId": product_id,
|
||||
"stolenFromTransport": stolen,
|
||||
"stored": stored,
|
||||
"lostDueToStorage": lost_storage,
|
||||
"branchId": branch_id
|
||||
},
|
||||
"lastOutcome": outcome
|
||||
});
|
||||
|
||||
save_result(pool, underground_id, &result)?;
|
||||
|
||||
if victim_app_uid > 0 {
|
||||
publish_transport_raid(broker, victim_app_uid, "transport_raided");
|
||||
publish_status(broker, victim_app_uid);
|
||||
publish_branch(broker, victim_app_uid);
|
||||
}
|
||||
|
||||
let ug_reason = match outcome {
|
||||
"major_success" => "raid_success",
|
||||
_ => "raid_partial_success",
|
||||
};
|
||||
publish_underground(broker, performer_app_uid, ug_reason);
|
||||
if stored > 0 {
|
||||
publish_underground(broker, performer_app_uid, "raid_loot_stored");
|
||||
}
|
||||
publish_status(broker, performer_app_uid);
|
||||
publish_branch(broker, performer_app_uid);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn finish_invalid(
|
||||
pool: &ConnectionPool,
|
||||
broker: &MessageBroker,
|
||||
underground_id: i32,
|
||||
performer_character_id: i32,
|
||||
code: &str,
|
||||
detail: Json,
|
||||
) -> Result<(), DbError> {
|
||||
let performer_app_uid = falukant_user_row(pool, performer_character_id)
|
||||
.map(|(_, app)| app)
|
||||
.unwrap_or(-1);
|
||||
|
||||
let result = json!({
|
||||
"status": "error",
|
||||
"code": code,
|
||||
"detail": detail
|
||||
});
|
||||
save_result(pool, underground_id, &result)?;
|
||||
if performer_app_uid > 0 {
|
||||
publish_underground(broker, performer_app_uid, "raid_repelled");
|
||||
publish_status(broker, performer_app_uid);
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn falukant_user_row(pool: &ConnectionPool, character_id: i32) -> Result<(i32, i32), DbError> {
|
||||
let mut conn = pool
|
||||
.get()
|
||||
.map_err(|e| DbError::new(format!("DB-Verbindung fehlgeschlagen: {e}")))?;
|
||||
conn.prepare("fu", QUERY_RAID_FALUKANT_USER_FOR_CHARACTER)
|
||||
.map_err(|e| DbError::new(format!("prepare fu: {e}")))?;
|
||||
let rows = conn
|
||||
.execute("fu", &[&character_id])
|
||||
.map_err(|e| DbError::new(format!("exec fu: {e}")))?;
|
||||
rows.first()
|
||||
.map(|r| {
|
||||
(
|
||||
parse_i32(r, "falukant_user_id", -1),
|
||||
parse_i32(r, "app_user_id", -1),
|
||||
)
|
||||
})
|
||||
.ok_or_else(|| DbError::new("character not found"))
|
||||
}
|
||||
|
||||
fn app_user_for_falukant(pool: &ConnectionPool, falukant_user_id: i32) -> Result<i32, DbError> {
|
||||
let mut conn = pool
|
||||
.get()
|
||||
.map_err(|e| DbError::new(format!("DB-Verbindung fehlgeschlagen: {e}")))?;
|
||||
conn.prepare("au", QUERY_RAID_APP_USER_FOR_FALUKANT)
|
||||
.map_err(|e| DbError::new(format!("prepare au: {e}")))?;
|
||||
let rows = conn
|
||||
.execute("au", &[&falukant_user_id])
|
||||
.map_err(|e| DbError::new(format!("exec au: {e}")))?;
|
||||
Ok(rows
|
||||
.first()
|
||||
.map(|r| parse_i32(r, "user_id", -1))
|
||||
.unwrap_or(-1))
|
||||
}
|
||||
|
||||
fn nearest_branch(
|
||||
pool: &ConnectionPool,
|
||||
raid_region_id: i32,
|
||||
performer_falukant_id: i32,
|
||||
) -> Result<(i32, i32), DbError> {
|
||||
let mut conn = pool
|
||||
.get()
|
||||
.map_err(|e| DbError::new(format!("DB-Verbindung fehlgeschlagen: {e}")))?;
|
||||
conn.prepare("nb", QUERY_RAID_NEAREST_BRANCH_FOR_USER)
|
||||
.map_err(|e| DbError::new(format!("prepare nb: {e}")))?;
|
||||
let rows = conn
|
||||
.execute("nb", &[&raid_region_id, &performer_falukant_id])
|
||||
.map_err(|e| DbError::new(format!("exec nb: {e}")))?;
|
||||
let branch_id = rows
|
||||
.first()
|
||||
.map(|r| parse_i32(r, "branch_id", -1))
|
||||
.unwrap_or(-1);
|
||||
let reg = rows
|
||||
.first()
|
||||
.map(|r| parse_i32(r, "region_id", -1))
|
||||
.unwrap_or(-1);
|
||||
Ok((branch_id, reg))
|
||||
}
|
||||
|
||||
fn add_to_branch_inventory(
|
||||
pool: &ConnectionPool,
|
||||
branch_id: i32,
|
||||
product_id: i32,
|
||||
quantity: i32,
|
||||
) -> Result<i32, DbError> {
|
||||
use std::cmp::min;
|
||||
|
||||
let stocks = load_stocks(pool, branch_id)?;
|
||||
let mut remaining = quantity;
|
||||
for s in stocks {
|
||||
if remaining <= 0 {
|
||||
break;
|
||||
}
|
||||
let free = s.total_capacity - s.filled;
|
||||
if free <= 0 {
|
||||
continue;
|
||||
}
|
||||
let to_store = min(remaining, free);
|
||||
store_in_stock(pool, s.stock_id, product_id, to_store)?;
|
||||
remaining -= to_store;
|
||||
}
|
||||
Ok(remaining)
|
||||
}
|
||||
|
||||
struct StockInfo {
|
||||
stock_id: i32,
|
||||
total_capacity: i32,
|
||||
filled: i32,
|
||||
}
|
||||
|
||||
fn load_stocks(pool: &ConnectionPool, branch_id: i32) -> Result<Vec<StockInfo>, DbError> {
|
||||
let mut conn = pool
|
||||
.get()
|
||||
.map_err(|e| DbError::new(format!("DB-Verbindung fehlgeschlagen: {e}")))?;
|
||||
conn.prepare("gs", QUERY_GET_AVAILABLE_STOCKS)
|
||||
.map_err(|e| DbError::new(format!("prepare gs: {e}")))?;
|
||||
let rows = conn
|
||||
.execute("gs", &[&branch_id])
|
||||
.map_err(|e| DbError::new(format!("exec gs: {e}")))?;
|
||||
let mut out = Vec::new();
|
||||
for row in rows {
|
||||
let stock_id = parse_i32(&row, "id", -1);
|
||||
let total_capacity = parse_i32(&row, "total_capacity", 0);
|
||||
let filled = parse_i32(&row, "filled", 0);
|
||||
if stock_id >= 0 {
|
||||
out.push(StockInfo {
|
||||
stock_id,
|
||||
total_capacity,
|
||||
filled,
|
||||
});
|
||||
}
|
||||
}
|
||||
Ok(out)
|
||||
}
|
||||
|
||||
fn store_in_stock(pool: &ConnectionPool, stock_id: i32, product_id: i32, quantity: i32) -> Result<(), DbError> {
|
||||
let mut conn = pool
|
||||
.get()
|
||||
.map_err(|e| DbError::new(format!("DB-Verbindung fehlgeschlagen: {e}")))?;
|
||||
conn.prepare("ins", QUERY_INSERT_INVENTORY)
|
||||
.map_err(|e| DbError::new(format!("prepare ins: {e}")))?;
|
||||
let quality: i32 = 100;
|
||||
conn.execute("ins", &[&stock_id, &product_id, &quantity, &quality])
|
||||
.map_err(|e| DbError::new(format!("exec ins: {e}")))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn save_result(pool: &ConnectionPool, underground_id: i32, result: &Json) -> Result<(), DbError> {
|
||||
let s = result.to_string();
|
||||
let mut conn = pool
|
||||
.get()
|
||||
.map_err(|e| DbError::new(format!("DB-Verbindung fehlgeschlagen: {e}")))?;
|
||||
conn.prepare("ur", QUERY_RAID_UPDATE_UNDERGROUND_RESULT)
|
||||
.map_err(|e| DbError::new(format!("prepare ur: {e}")))?;
|
||||
conn.execute("ur", &[&underground_id, &s])
|
||||
.map_err(|e| DbError::new(format!("exec ur: {e}")))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn publish_transport_raid(broker: &MessageBroker, user_id: i32, reason: &str) {
|
||||
let m = format!(
|
||||
r#"{{"event":"falukantTransportRaid","user_id":{},"reason":"{}"}}"#,
|
||||
user_id, reason
|
||||
);
|
||||
broker.publish(m);
|
||||
}
|
||||
|
||||
fn publish_underground(broker: &MessageBroker, user_id: i32, reason: &str) {
|
||||
if user_id <= 0 {
|
||||
return;
|
||||
}
|
||||
let m = format!(
|
||||
r#"{{"event":"falukantUndergroundUpdate","user_id":{},"reason":"{}"}}"#,
|
||||
user_id, reason
|
||||
);
|
||||
broker.publish(m);
|
||||
}
|
||||
|
||||
fn publish_status(broker: &MessageBroker, user_id: i32) {
|
||||
if user_id <= 0 {
|
||||
return;
|
||||
}
|
||||
let m = format!(r#"{{"event":"falukantUpdateStatus","user_id":{}}}"#, user_id);
|
||||
broker.publish(m);
|
||||
}
|
||||
|
||||
fn publish_branch(broker: &MessageBroker, user_id: i32) {
|
||||
if user_id <= 0 {
|
||||
return;
|
||||
}
|
||||
let m = format!(r#"{{"event":"falukantBranchUpdate","user_id":{}}}"#, user_id);
|
||||
broker.publish(m);
|
||||
}
|
||||
|
||||
fn parse_i32(row: &Row, key: &str, default: i32) -> i32 {
|
||||
row.get(key)
|
||||
.and_then(|v| v.parse::<i32>().ok())
|
||||
.unwrap_or(default)
|
||||
}
|
||||
Reference in New Issue
Block a user