Skip to content

AbstractRTask

Base class for R-based federated learning tasks. Provides a Python-R bridge that executes R scripts via subprocess and communicates through JSON files.

starfish.controller.tasks.abstract_r_task.AbstractRTask

Bases: AbstractTask, ABC

Base class for FL tasks whose core logic is written in R.

Subclasses must set r_script_dir to the directory containing prepare_data.R, training.R, and aggregate.R.

Source code in controller/starfish/controller/tasks/abstract_r_task.py
 16
 17
 18
 19
 20
 21
 22
 23
 24
 25
 26
 27
 28
 29
 30
 31
 32
 33
 34
 35
 36
 37
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
class AbstractRTask(AbstractTask, ABC):
    """
    Base class for FL tasks whose core logic is written in R.

    Subclasses must set ``r_script_dir`` to the directory containing
    ``prepare_data.R``, ``training.R``, and ``aggregate.R``.
    """
    r_script_dir: str = None  # set by subclass

    def __init__(self, run):
        super().__init__(run)
        self.sample_size = None

    # ------------------------------------------------------------------
    # R bridge helpers
    # ------------------------------------------------------------------

    def _run_r_script(self, script_name, input_json_path, output_json_path):
        """Run an R script via ``Rscript`` and return True on success."""
        script_path = os.path.join(self.r_script_dir, script_name)
        cmd = ['Rscript', '--vanilla', script_path,
               input_json_path, output_json_path]
        self.logger.debug("Running R script: {}".format(' '.join(cmd)))
        result = subprocess.run(
            cmd, capture_output=True, text=True, timeout=600)
        if result.stdout:
            self.logger.debug("R stdout: {}".format(result.stdout))
        if result.stderr:
            self.logger.debug("R stderr: {}".format(result.stderr))
        if result.returncode != 0:
            self.logger.warning(
                "R script {} failed with exit code {}: {}".format(
                    script_name, result.returncode,
                    result.stderr.strip() if result.stderr else '(no stderr)'))
            return False
        return True

    def _write_input_json(self, extra=None):
        """Write input JSON for an R script and return the file path."""
        data = {
            'run_id': self.run_id,
            'project_id': self.project_id,
            'batch_id': self.batch_id,
            'cur_seq': self.cur_seq,
            'round': self.get_round(),
            'config': self.tasks[self.cur_seq - 1].get('config', {}),
        }
        if extra:
            data.update(extra)
        fd, path = tempfile.mkstemp(suffix='.json', prefix='r_input_')
        with os.fdopen(fd, 'w') as f:
            json.dump(data, f)
        return path

    def _make_output_path(self):
        fd, path = tempfile.mkstemp(suffix='.json', prefix='r_output_')
        os.close(fd)
        return path

    def _read_output_json(self, path):
        with open(path, 'r') as f:
            return json.load(f)

    # ------------------------------------------------------------------
    # AbstractTask interface
    # ------------------------------------------------------------------

    def validate(self) -> bool:
        task_round = self.get_round()
        self.logger.debug(
            "Run {} - task {} - round {} task begins".format(
                self.run_id, self.cur_seq, task_round))
        return self.download_artifact()

    def prepare_data(self) -> bool:
        self.logger.debug(
            'Loading dataset for run {} ...'.format(self.run_id))

        dataset_url = gen_dataset_url(self.run_id)
        data_path = dataset_url + 'dataset' if dataset_url else None
        if not data_path or not os.path.exists(data_path):
            self.logger.warning("Dataset not found at {}".format(data_path))
            return False

        extra = {'data_path': data_path}

        # Load previous model if not first round
        if not self.is_first_round():
            previous_model = self._load_previous_model()
            if previous_model:
                extra['previous_model'] = previous_model

        input_path = self._write_input_json(extra)
        output_path = self._make_output_path()

        try:
            if not self._run_r_script('prepare_data.R', input_path, output_path):
                return False
            result = self._read_output_json(output_path)
            self.sample_size = result.get('sample_size', 0)
            return result.get('valid', False)
        finally:
            self._cleanup_temp(input_path, output_path)

    def training(self) -> bool:
        self.logger.info('Starting R training...')

        dataset_url = gen_dataset_url(self.run_id)
        data_path = dataset_url + 'dataset' if dataset_url else None

        extra = {'data_path': data_path, 'sample_size': self.sample_size}

        if not self.is_first_round():
            previous_model = self._load_previous_model()
            if previous_model:
                extra['previous_model'] = previous_model

        input_path = self._write_input_json(extra)
        output_path = self._make_output_path()

        try:
            if not self._run_r_script('training.R', input_path, output_path):
                return False
            result = self._read_output_json(output_path)
            self.logger.info('R training complete.')
            url = gen_mid_artifacts_url(
                self.run_id, self.cur_seq, self.get_round())
            self.logger.info(
                "Upload: {} \n to: {}".format(result, url))
            return self.save_artifacts(url, json.dumps(result))
        finally:
            self._cleanup_temp(input_path, output_path)

    def do_aggregate(self) -> bool:
        mid_artifacts = self._collect_mid_artifacts()
        if not mid_artifacts:
            self.logger.warning(
                "No mid-artifacts found for aggregation")
            return False

        extra = {'mid_artifacts': mid_artifacts}
        input_path = self._write_input_json(extra)
        output_path = self._make_output_path()

        try:
            if not self._run_r_script('aggregate.R', input_path, output_path):
                return False
            result = self._read_output_json(output_path)
            self.sample_size = result.get('sample_size', 0)
            url = gen_artifacts_url(
                self.run_id, self.cur_seq, self.get_round())
            self.logger.info(
                "Upload: {} \n to: {}".format(result, url))
            if self.save_artifacts(url, json.dumps(result)):
                self.upload(True)
                return True
            return False
        finally:
            self._cleanup_temp(input_path, output_path)

    # ------------------------------------------------------------------
    # Private helpers
    # ------------------------------------------------------------------

    def _load_previous_model(self):
        seq_no, round_no = self.get_previous_seq_and_round()
        if not seq_no or not round_no:
            return None
        directory = downloaded_artifacts_url(self.run_id, seq_no, round_no)
        if not directory:
            return None
        for path in Path(directory).rglob(
                "*-{}-{}-artifacts".format(seq_no, round_no)):
            with open(str(path), 'r') as f:
                for line in f:
                    return json.loads(line)
        return None

    def _collect_mid_artifacts(self):
        artifacts = []
        directory = gen_all_mid_artifacts_url(self.project_id, self.batch_id)
        if not directory:
            return artifacts
        for path in Path(directory).rglob(
                "*-{}-{}-mid-artifacts".format(self.cur_seq, self.get_round())):
            with open(str(path), 'r') as f:
                for line in f:
                    artifacts.append(json.loads(line))
        return artifacts

    def _cleanup_temp(self, *paths):
        for p in paths:
            try:
                os.unlink(p)
            except OSError:
                pass