Add KSEI archive fallback import

This commit is contained in:
Rasyidan Akbar F. 2026-03-31 19:29:00 +07:00
commit 1f5c4b2798
11 changed files with 622 additions and 59 deletions

View file

@ -18,7 +18,7 @@ use crate::output::table::format_idr;
use crate::ownership::types::{
ChangeType, FlowSignal, HolderRow, KseiHolding, OwnershipRelease, OwnershipSource,
};
use crate::ownership::{db, entities, graph, parser, remote, search, snapshot};
use crate::ownership::{archive, db, entities, graph, parser, remote, search, snapshot};
#[derive(Debug, Args)]
pub struct OwnershipCmd {
@ -30,7 +30,7 @@ pub struct OwnershipCmd {
pub enum OwnershipCommand {
/// Discover the latest IDX-hosted ownership report URLs.
Discover(DiscoverArgs),
/// Import ownership data from KSEI PDF or Bing API.
/// Import ownership data from KSEI PDF or archive fallback files.
Import(ImportArgs),
/// Install or refresh a maintained ownership SQLite snapshot.
Sync(SyncArgs),
@ -71,7 +71,7 @@ pub struct ImportArgs {
/// URL to a remote ownership PDF.
#[arg(long)]
pub url: Option<String>,
/// Path to local KSEI PDF file.
/// Path to local KSEI ownership PDF, ZIP, or TXT file.
#[arg(long)]
pub file: Option<PathBuf>,
/// Fetch Bing institutional data for these symbols.
@ -773,30 +773,40 @@ fn handle_import(args: &ImportArgs, config: &IdxConfig) -> Result<(), IdxError>
}
}
let Some(pdf_input) = resolve_pdf_input(args)? else {
let Some(import_input) = resolve_import_input(args)? else {
return Ok(());
};
let conn = db::open_db(config)?;
let sha256 = sha256_file(&pdf_input.pdf_path)?;
let sha256 = sha256_file(&import_input.import_path)?;
if !args.force && db::release_exists(&conn, &sha256)? {
println!("Release already imported (sha256: {sha256}). Use --force to re-import.");
return Ok(());
}
let raw_rows = parser::parse_ksei_pdf(&pdf_input.pdf_path)?;
if raw_rows.is_empty() {
return Err(IdxError::ParseError(
"no KSEI rows parsed from PDF".to_string(),
));
}
let drafts = match import_input.format {
ImportInputFormat::Pdf => {
let raw_rows = parser::parse_ksei_pdf(&import_input.import_path)?;
if raw_rows.is_empty() {
return Err(IdxError::ParseError(
"no KSEI rows parsed from PDF".to_string(),
));
}
let mut holdings = Vec::with_capacity(raw_rows.len());
let mut drafts = Vec::with_capacity(raw_rows.len());
for raw in &raw_rows {
drafts.push(entities::normalize_ksei_row(raw)?);
}
drafts
}
ImportInputFormat::Archive => archive::parse_balancepos_file(&import_input.import_path)?,
};
let mut holdings = Vec::with_capacity(drafts.len());
let mut ticker_ids = HashSet::new();
for raw in &raw_rows {
let draft = entities::normalize_ksei_row(raw)?;
for draft in drafts {
let ticker_id = db::upsert_ticker(&conn, &draft.ticker_code, draft.issuer_name.as_deref())?;
let entity_id =
entities::resolve_entity(&conn, &draft.raw_investor_name, OwnershipSource::Ksei)?;
@ -831,7 +841,7 @@ fn handle_import(args: &ImportArgs, config: &IdxConfig) -> Result<(), IdxError>
let release = OwnershipRelease {
id: 0,
source_url: pdf_input.source_url,
source_url: import_input.source_url,
sha256,
as_of_date,
row_count: inserted_rows,
@ -900,23 +910,31 @@ fn parse_discovery_family(raw: &str) -> Result<Option<remote::OwnershipReportFam
}
}
struct ResolvedPdfInput {
pdf_path: PathBuf,
struct ResolvedImportInput {
import_path: PathBuf,
source_url: Option<String>,
format: ImportInputFormat,
}
fn resolve_pdf_input(args: &ImportArgs) -> Result<Option<ResolvedPdfInput>, IdxError> {
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ImportInputFormat {
Pdf,
Archive,
}
fn resolve_import_input(args: &ImportArgs) -> Result<Option<ResolvedImportInput>, IdxError> {
if let Some(path) = &args.file {
if !path.exists() {
return Err(IdxError::Io(format!(
"input PDF not found: {}",
"input ownership file not found: {}",
path.display()
)));
}
return Ok(Some(ResolvedPdfInput {
pdf_path: path.clone(),
return Ok(Some(ResolvedImportInput {
import_path: path.clone(),
source_url: None,
format: detect_local_import_format(path)?,
}));
}
@ -925,15 +943,34 @@ fn resolve_pdf_input(args: &ImportArgs) -> Result<Option<ResolvedPdfInput>, IdxE
validate_import_url(trimmed)?;
let target = cache_pdf_path(trimmed)?;
download_pdf(trimmed, &target)?;
return Ok(Some(ResolvedPdfInput {
pdf_path: target,
return Ok(Some(ResolvedImportInput {
import_path: target,
source_url: Some(trimmed.to_string()),
format: ImportInputFormat::Pdf,
}));
}
Ok(None)
}
fn detect_local_import_format(path: &Path) -> Result<ImportInputFormat, IdxError> {
match path
.extension()
.and_then(|value| value.to_str())
.map(|value| value.trim().to_ascii_lowercase())
.as_deref()
{
Some("pdf") => Ok(ImportInputFormat::Pdf),
Some("zip") | Some("txt") if archive::supports_local_archive_file(path) => {
Ok(ImportInputFormat::Archive)
}
_ => Err(IdxError::InvalidInput(format!(
"unsupported local ownership file {}; expected a .pdf, .zip, or .txt input",
path.display()
))),
}
}
fn cache_pdf_path(url: &str) -> Result<PathBuf, IdxError> {
let dirs = ProjectDirs::from("", "", "idx")
.ok_or_else(|| IdxError::Io("unable to resolve cache directory".to_string()))?;

360
src/ownership/archive.rs Normal file
View file

@ -0,0 +1,360 @@
use std::fs;
use std::io::{Cursor, Read};
use std::path::Path;
use chrono::NaiveDate;
use crate::error::IdxError;
use crate::ownership::types::{InvestorTypeCode, KseiHoldingDraft, Locality};
const EXPECTED_HEADER: &[&str] = &[
"Date",
"Code",
"Type",
"Sec. Num",
"Price",
"Local IS",
"Local CP",
"Local PF",
"Local IB",
"Local ID",
"Local MF",
"Local SC",
"Local FD",
"Local OT",
"Total",
"Foreign IS",
"Foreign CP",
"Foreign PF",
"Foreign IB",
"Foreign ID",
"Foreign MF",
"Foreign SC",
"Foreign FD",
"Foreign OT",
"Total",
];
const LOCAL_BUCKETS: &[(&str, usize)] = &[
("IS", 5),
("CP", 6),
("PF", 7),
("IB", 8),
("ID", 9),
("MF", 10),
("SC", 11),
("FD", 12),
("OT", 13),
];
const FOREIGN_BUCKETS: &[(&str, usize)] = &[
("IS", 15),
("CP", 16),
("PF", 17),
("IB", 18),
("ID", 19),
("MF", 20),
("SC", 21),
("FD", 22),
("OT", 23),
];
pub fn supports_local_archive_file(path: &Path) -> bool {
path.extension()
.and_then(|value| value.to_str())
.map(|value| {
let ext = value.trim().to_ascii_lowercase();
ext == "txt" || ext == "zip"
})
.unwrap_or(false)
}
pub fn parse_balancepos_file(path: &Path) -> Result<Vec<KseiHoldingDraft>, IdxError> {
let raw = match path
.extension()
.and_then(|value| value.to_str())
.map(|value| value.trim().to_ascii_lowercase())
.as_deref()
{
Some("txt") => fs::read_to_string(path).map_err(|e| {
IdxError::Io(format!(
"failed to read KSEI archive TXT {}: {e}",
path.display()
))
})?,
Some("zip") => extract_txt_from_zip(path)?,
_ => {
return Err(IdxError::InvalidInput(format!(
"unsupported local ownership archive file {}; expected .txt or .zip",
path.display()
)));
}
};
parse_balancepos_text(&raw)
}
pub fn parse_balancepos_text(raw: &str) -> Result<Vec<KseiHoldingDraft>, IdxError> {
let mut lines = raw.lines().filter(|line| !line.trim().is_empty());
let header = lines
.next()
.ok_or_else(|| IdxError::ParseError("empty KSEI archive TXT input".to_string()))?;
validate_header(header)?;
let mut drafts = Vec::new();
for (line_number, line) in lines.enumerate() {
let columns: Vec<&str> = line.split('|').map(str::trim).collect();
if columns.len() != EXPECTED_HEADER.len() {
return Err(IdxError::ParseError(format!(
"invalid KSEI archive TXT row {}: expected {} columns, got {}",
line_number + 2,
EXPECTED_HEADER.len(),
columns.len()
)));
}
if !columns[2].eq_ignore_ascii_case("EQUITY") {
continue;
}
let report_date = parse_archive_date(columns[0])?;
let ticker_code = columns[1].trim().to_uppercase();
let sec_num = parse_archive_number(columns[3], line_number + 2, "Sec. Num")?;
if sec_num <= 0 {
continue;
}
append_bucket_drafts(
&mut drafts,
&ticker_code,
report_date,
sec_num,
columns.as_slice(),
Locality::Local,
LOCAL_BUCKETS,
)?;
append_bucket_drafts(
&mut drafts,
&ticker_code,
report_date,
sec_num,
columns.as_slice(),
Locality::Foreign,
FOREIGN_BUCKETS,
)?;
}
if drafts.is_empty() {
return Err(IdxError::ParseError(
"no importable EQUITY rows found in KSEI archive TXT".to_string(),
));
}
Ok(drafts)
}
fn extract_txt_from_zip(path: &Path) -> Result<String, IdxError> {
let bytes = fs::read(path).map_err(|e| {
IdxError::Io(format!(
"failed to read KSEI archive ZIP {}: {e}",
path.display()
))
})?;
let cursor = Cursor::new(bytes);
let mut zip = zip::ZipArchive::new(cursor)
.map_err(|e| IdxError::ParseError(format!("failed to open KSEI archive ZIP: {e}")))?;
for index in 0..zip.len() {
let mut file = zip.by_index(index).map_err(|e| {
IdxError::ParseError(format!("failed to read KSEI archive ZIP entry: {e}"))
})?;
if file.is_dir() {
continue;
}
let name = file.name().to_ascii_lowercase();
if !name.ends_with(".txt") {
continue;
}
let mut output = String::new();
file.read_to_string(&mut output).map_err(|e| {
IdxError::ParseError(format!("failed to decode KSEI archive TXT entry: {e}"))
})?;
return Ok(output);
}
Err(IdxError::ParseError(
"KSEI archive ZIP did not contain a TXT payload".to_string(),
))
}
fn validate_header(header: &str) -> Result<(), IdxError> {
let columns: Vec<&str> = header.split('|').map(str::trim).collect();
if columns != EXPECTED_HEADER {
return Err(IdxError::ParseError(
"KSEI archive TXT header did not match the expected balancepos layout".to_string(),
));
}
Ok(())
}
fn parse_archive_date(raw: &str) -> Result<NaiveDate, IdxError> {
let trimmed = raw.trim();
if trimmed.len() != 11 {
return Err(IdxError::ParseError(format!(
"invalid KSEI archive date `{trimmed}`"
)));
}
let canonical = format!(
"{}-{}-{}",
&trimmed[..2],
titlecase_month(&trimmed[3..6]),
&trimmed[7..11]
);
NaiveDate::parse_from_str(&canonical, "%d-%b-%Y")
.map_err(|e| IdxError::ParseError(format!("invalid KSEI archive date `{trimmed}`: {e}")))
}
fn titlecase_month(raw: &str) -> String {
let upper = raw.trim().to_ascii_uppercase();
let mut chars = upper.chars();
match chars.next() {
Some(first) => {
let mut output = String::new();
output.push(first.to_ascii_uppercase());
output.push_str(&chars.as_str().to_ascii_lowercase());
output
}
None => String::new(),
}
}
fn parse_archive_number(raw: &str, line_number: usize, field: &str) -> Result<i64, IdxError> {
raw.trim().parse::<i64>().map_err(|e| {
IdxError::ParseError(format!(
"invalid KSEI archive TXT value in row {line_number} field `{field}`: {e}"
))
})
}
fn append_bucket_drafts(
drafts: &mut Vec<KseiHoldingDraft>,
ticker_code: &str,
report_date: NaiveDate,
sec_num: i64,
columns: &[&str],
locality: Locality,
buckets: &[(&str, usize)],
) -> Result<(), IdxError> {
for (investor_type, column_index) in buckets {
let shares =
parse_archive_number(columns[*column_index], 0, investor_type).map_err(|_| {
IdxError::ParseError(format!(
"invalid KSEI archive TXT share count for {ticker_code} {investor_type}"
))
})?;
if shares <= 0 {
continue;
}
drafts.push(KseiHoldingDraft {
ticker_code: ticker_code.to_string(),
issuer_name: None,
raw_investor_name: synthetic_holder_name(locality, investor_type),
investor_type: Some(InvestorTypeCode((*investor_type).to_string())),
locality: Some(locality),
nationality: None,
domicile: None,
holdings_scripless: shares,
holdings_scrip: 0,
total_shares: shares,
percentage_bps: compute_percentage_bps(shares, sec_num),
report_date,
});
}
Ok(())
}
fn synthetic_holder_name(locality: Locality, investor_type: &str) -> String {
let prefix = match locality {
Locality::Local => "LOCAL",
Locality::Foreign => "FOREIGN",
};
format!("KSEI AGGREGATE {prefix} {investor_type}")
}
fn compute_percentage_bps(shares: i64, sec_num: i64) -> i64 {
let shares_i128 = i128::from(shares);
let sec_num_i128 = i128::from(sec_num);
let rounded = ((shares_i128 * 10_000) + (sec_num_i128 / 2)) / sec_num_i128;
i64::try_from(rounded).unwrap_or(0)
}
#[cfg(test)]
mod tests {
use super::parse_balancepos_text;
use crate::ownership::entities::normalize_ksei_row;
use crate::ownership::parser::parse_stext_xml;
use crate::ownership::types::Locality;
#[test]
fn parses_balancepos_excerpt_into_bucket_holders() {
let raw = include_str!("../../tests/fixtures/ksei_balancepos_20260227_excerpt.txt");
let drafts = parse_balancepos_text(raw).expect("balancepos excerpt parses");
let local_cp = drafts
.iter()
.find(|draft| {
draft.ticker_code == "AADI"
&& draft.investor_type.as_ref().map(|code| code.0.as_str()) == Some("CP")
&& draft.locality == Some(Locality::Local)
})
.expect("AADI local CP bucket");
assert_eq!(local_cp.raw_investor_name, "KSEI AGGREGATE LOCAL CP");
assert_eq!(local_cp.total_shares, 5_035_745_466);
assert_eq!(local_cp.percentage_bps, 6467);
}
#[test]
fn balancepos_cross_check_contains_pdf_holder_bucket() {
let pdf_rows = parse_stext_xml(include_str!(
"../../tests/fixtures/ksei_above1_stext_excerpt.xml"
))
.expect("pdf fixture rows");
let pdf_drafts: Vec<_> = pdf_rows
.iter()
.map(normalize_ksei_row)
.collect::<Result<_, _>>()
.expect("normalized pdf drafts");
let archive_drafts = parse_balancepos_text(include_str!(
"../../tests/fixtures/ksei_balancepos_20260227_excerpt.txt"
))
.expect("archive drafts");
let pdf_local_cp = pdf_drafts
.iter()
.find(|draft| {
draft.ticker_code == "AADI"
&& draft.investor_type.as_ref().map(|code| code.0.as_str()) == Some("CP")
&& draft.locality == Some(Locality::Local)
})
.expect("pdf local cp");
let archive_local_cp = archive_drafts
.iter()
.find(|draft| {
draft.ticker_code == "AADI"
&& draft.investor_type.as_ref().map(|code| code.0.as_str()) == Some("CP")
&& draft.locality == Some(Locality::Local)
})
.expect("archive local cp");
assert_eq!(pdf_local_cp.report_date, archive_local_cp.report_date);
assert!(pdf_local_cp.total_shares <= archive_local_cp.total_shares);
assert!(pdf_local_cp.percentage_bps <= archive_local_cp.percentage_bps);
}
}

View file

@ -1,3 +1,4 @@
pub mod archive;
pub mod db;
pub mod entities;
pub mod graph;