"""Contract creation processor — business logic orchestrator.""" import logging from connectors.abacus import AbacusClient from domain.contract_input_builder import ContractInputBuilder from infra.output import NullOutputSink, OutputSink from infra.resume import load_processed_identifiers from infra.row_reader import RowReaderProtocol, RowValidationError from schemas import ( DRY_RUN_CONTRACT_ID, AbacusContract, AbacusContractWithLifecyclesInput, BuildSkip, ContractRow, ErrorRow, GraphQLResult, MissingFieldPolicy, ProcessingResults, ProcessingStatus, ResultColumn, SkippedRow, SkipReason, ) logger = logging.getLogger(__name__) def _make_identifier(row: ContractRow) -> tuple[str, str] | None: """Build a (account_id, contract_name) identifier from a ContractRow. Returns None when either component is empty. """ if row.account_id and row.contract_name: return (str(row.account_id), row.contract_name) return None class ContractProcessor: """Orchestrates batch contract creation from CSV. All services are injected via the constructor. The processor does not know about CLI concerns, logging configuration, or process lifecycle. """ def __init__( self, abacus: AbacusClient, contract_builder: ContractInputBuilder, row_reader: RowReaderProtocol, dry_run: bool = False, output_sink: OutputSink | None = None, ): self._abacus = abacus self._contract_builder = contract_builder self._row_reader = row_reader self._dry_run = dry_run self._output: OutputSink = output_sink or NullOutputSink() def process( self, input_file: str, limit: int | None = None, policy: MissingFieldPolicy = MissingFieldPolicy.DEFAULT, resume_csv: str | None = None, ) -> ProcessingResults: """Process CSV file and create contracts. Args: input_file: Path to input file (CSV, JSON, or XLSX). limit: Maximum number of rows to process. policy: How to handle missing optional fields. resume_csv: Path to existing success CSV for resume state. Raises ApiError / GraphQLError on API failures. Raises FileNotFoundError if the input file doesn't exist. Raises MissingColumnsError if required columns are absent. """ results = ProcessingResults() processed = load_processed_identifiers(resume_csv) seen: set[tuple[str, str]] = set() row_iter = self._row_reader.read(input_file) self._output.open(row_iter.fieldnames) rows_processed = 0 for idx, row_or_error in row_iter: if limit and rows_processed >= limit: logger.info(f'Reached limit of {limit} rows') break rows_processed += 1 if isinstance(row_or_error, RowValidationError): results.skipped.append( SkippedRow( row=idx, data=row_or_error.raw_data, reason=row_or_error.error, ) ) continue row = row_or_error result, skip_reason = self._process_row( row, idx, policy, processed, seen ) if skip_reason: if skip_reason != SkipReason.ALREADY_PROCESSED: results.skipped.append( SkippedRow( row=idx, data=row.model_dump(by_alias=True), reason=skip_reason, ) ) continue if result: if result.status in ( ProcessingStatus.SUCCESS, ProcessingStatus.DRY_RUN, ): results.success.append(result) success_row_data = row.model_dump(by_alias=True) if result.status == ProcessingStatus.SUCCESS and result.data: success_row_data[ResultColumn.CONTRACT_ID] = ( result.data.contract_id ) else: success_row_data[ResultColumn.CONTRACT_ID] = DRY_RUN_CONTRACT_ID self._output.add_success(success_row_data) else: results.errors.append( ErrorRow( row=idx, data=row.model_dump(by_alias=True), result=result, ) ) if results.errors or results.skipped: self._output.write_failures(results.errors, results.skipped) self._output.close() return results # ------------------------------------------------------------------ # Private helpers # ------------------------------------------------------------------ def _process_row( self, row: ContractRow, row_idx: int, policy: MissingFieldPolicy, processed: set[tuple[str, str]], seen: set[tuple[str, str]], ) -> tuple[GraphQLResult[AbacusContract] | None, str | None]: """Validate, build payload, and create contract from a typed row.""" identifier = _make_identifier(row) if identifier is not None and identifier in seen: logger.warning( f'[DUPLICATE] Row {row_idx} is duplicate in input file: ' f'{row.contract_name} (Account: {row.account_id})' ) return None, SkipReason.DUPLICATE_IN_INPUT if identifier is not None and identifier in processed: logger.info( f'[SKIP] Row {row_idx} already processed: ' f'{row.contract_name} (Account: {row.account_id})' ) return None, SkipReason.ALREADY_PROCESSED if identifier is not None: seen.add(identifier) logger.info( f'Processing row {row_idx}: {row.contract_name} (Account: {row.account_id})' ) build_result = self._contract_builder.build(row, policy) if isinstance(build_result, BuildSkip): logger.warning(f'Row {row_idx}: skipped ({build_result.reason})') return None, SkipReason.INVALID_DATA if build_result.defaults_applied: logger.info( f'Row {row_idx}: defaults applied: ' f'{", ".join(build_result.defaults_applied)}' ) result = self._create_and_attach(build_result.input) if ( result.status in (ProcessingStatus.SUCCESS, ProcessingStatus.DRY_RUN) and identifier is not None ): processed.add(identifier) return result, None def _create_and_attach( self, gql_input: AbacusContractWithLifecyclesInput, ) -> GraphQLResult[AbacusContract]: """Create a contract with lifecycle and run controller. The GraphQL mutation accepts runControllerId as part of the contract input, so creation and attachment happen in one call. """ if self._dry_run: logger.info( f'[DRY RUN] Would create contract: {gql_input.contract.contract_name}' ) if gql_input.contract.run_controller_id is not None: logger.info( f'[DRY RUN] run_controller_id: ' f'{gql_input.contract.run_controller_id}, ' f'signing_entity_id: ' f'{gql_input.contract.reference_signing_entity_id}' ) return GraphQLResult(status=ProcessingStatus.DRY_RUN) return self._abacus.create_contract_with_lifecycles(gql_input)